7.9 Worked examples
① Why mod N is so bad. With N nodes and a key space of K keys, moving from N to N+1 nodes:
mod N: a key stays put only ifh mod N == h mod (N+1)— roughlyK/(N+1)keys stay, so ~N/(N+1)of ALL keys move. Going 3→4 nodes moves ~75% of the data.- Fixed shards / consistent hashing: only
K/(N+1)keys move — going 3→4 moves ~25%, which is the theoretical minimum (the new node's fair share).
② Choosing a fixed shard count. You expect to grow from 3 nodes to at most 60 nodes.
- Minimum shards = 60 (you can't have more nodes than shards)
- Pick a highly divisible number: 720 = 2⁴·3²·5, divisible by 1,2,3,4,5,6,8,9,10,12,15,16,18,20,24,30,36,40,45,48,60… so it splits evenly across almost any node count.
- At 3 nodes: 240 shards/node. At 60 nodes: 12 shards/node.
- If total data reaches 7.2 TB, each shard is 10 GB — in the "just right" band. At 72 TB each shard is 100 GB and rebalancing/recovery become expensive — that's your resharding trigger.
③ Hot key salting arithmetic. A celebrity key takes 50,000 writes/s and 200,000 reads/s; a shard handles 10,000 ops/s.
- Unsalted: one shard sees 250,000 ops/s → 25× over capacity.
- Salt with 2 digits (100 keys): writes become 500/s per salted key across up to 100 shards ✔. But every read must query all 100 keys, so total read work becomes 200,000 × 100 = 20,000,000 key-reads/s — catastrophically worse.
- Correct design: salt writes, and serve reads from a cache or a materialized aggregate, not by gathering 100 keys. This is why the book says salting splits only the write load.
④ Scatter-query cost. 100 shards, each shard's p99 = 10 ms, p50 = 2 ms. A scatter query waits for the slowest shard.
- P(all 100 shards respond under their p99) = 0.99¹⁰⁰ ≈ 36%
- So ~64% of scatter queries take ≥10 ms, i.e. the cluster's p50 for scatter queries ≈ each shard's p99. That is tail latency amplification, and it's why local secondary indexes limit scalability.
⑤ Global index write amplification. A document with 40 indexed terms, global index sharded over 20 shards. One document write may touch up to min(40, 20) = 20 index shards plus 1 data shard. Either you run a 21-participant distributed transaction (slow, Ch 8) or you accept asynchronous, stale indexes (DynamoDB's choice).
⑥ Partition size and bucketing. IoT: 10,000 sensors × 1 reading/s × 200 bytes.
- Partition key
sensor_id: each partition grows 17.3 MB/day → 6.3 GB/year — over the 10 GB "big partition" threshold in under two years, and unbounded thereafter. - Partition key
(sensor_id, yyyy-mm): each partition caps at ~520 MB ✔, andWHERE sensor_id = X AND ts BETWEEN …still hits one partition per month. - Cost: a query spanning a year now touches 12 partitions instead of 1 — the deliberate trade.