Learn Labs
7. Sharding

7.7 Production failure catalog for this chapter

Production failure catalog
0 rows
SymptomUnderlying mechanism
9 of 10 nodes idle, one at 100%Hot shard — skewed partition key
All writes hit one shard, others idleTimestamp as the first element of a key-range partition key
Adding a node moves nearly all the datahash mod N
Can't add more nodes than shardsFixed shard count chosen too low
Resharding requires downtimeSystem doesn't allow resharding with concurrent writes
Shard split makes an already-hot shard worseSplitting rewrites all the data, and the shard needing a split is usually the loaded one
Cluster death spiral after one node slowsAuto-rebalance + auto-failure-detection cascading failure
Rebalancing can't keep up with writesSystem near max write throughput while splitting
A single celebrity melts one nodeHot key — uniform hashing ≠ uniform load
Salting the hot key didn't help readsSalting splits writes only; reads must still gather all 100 keys
One query silently costs 100×Scatter query — partition key omitted
More shards, same query throughputLocal secondary index — every shard processes every query
Secondary index disagrees with the dataHand-rolled index + race conditions / partial write failures
Global index returns stale resultsAsynchronous global index maintenance (DynamoDB)
Two coordinators assign the same shard differentlySplit brain in the shard-assignment coordinator
Requests lost during a shard moveCutover window with in-flight requests to the old node
One partition grows to 8 GB and everything times outUnbounded partition — missing bucketing in the clustering key
Cross-shard write half-succeededNo distributed transaction (Ch 8)