7. Sharding
7.7 Production failure catalog for this chapter
| Symptom | Underlying mechanism |
|---|---|
| 9 of 10 nodes idle, one at 100% | Hot shard — skewed partition key |
| All writes hit one shard, others idle | Timestamp as the first element of a key-range partition key |
| Adding a node moves nearly all the data | hash mod N |
| Can't add more nodes than shards | Fixed shard count chosen too low |
| Resharding requires downtime | System doesn't allow resharding with concurrent writes |
| Shard split makes an already-hot shard worse | Splitting rewrites all the data, and the shard needing a split is usually the loaded one |
| Cluster death spiral after one node slows | Auto-rebalance + auto-failure-detection cascading failure |
| Rebalancing can't keep up with writes | System near max write throughput while splitting |
| A single celebrity melts one node | Hot key — uniform hashing ≠ uniform load |
| Salting the hot key didn't help reads | Salting splits writes only; reads must still gather all 100 keys |
| One query silently costs 100× | Scatter query — partition key omitted |
| More shards, same query throughput | Local secondary index — every shard processes every query |
| Secondary index disagrees with the data | Hand-rolled index + race conditions / partial write failures |
| Global index returns stale results | Asynchronous global index maintenance (DynamoDB) |
| Two coordinators assign the same shard differently | Split brain in the shard-assignment coordinator |
| Requests lost during a shard move | Cutover window with in-flight requests to the old node |
| One partition grows to 8 GB and everything times out | Unbounded partition — missing bucketing in the clustering key |
| Cross-shard write half-succeeded | No distributed transaction (Ch 8) |