Sharding, Hot Keys & Rebalancing
The Product Page That Melted Redis
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
Part 0. Start here
The problem: 66,667 reads and 10,000 likes a second on one post
A celebrity with 60 million followers posts a photo. The push notification goes out, and about 8 million followers open the app over the next 2 minutes. Every feed read pulls her newest posts with one Query on one DynamoDB partition key. The post's like counter is a single item, and it takes about 10,000 likes a second (our example's number). The table has plenty of capacity: on average it is about 20% busy (also our example's number).
The arithmetic, with DynamoDB's documented limits:
| What | Arithmetic | Result |
|---|---|---|
| Reads on one key | 8,000,000 reads ÷ 120 s | 66,667 a second |
| What one partition serves | At most 3,000 read units and 1,000 write units a second | a fixed ceiling |
| One feed query | 20 posts ≈ 11 KB = three 4 KB units = 1.5 read units when eventually consistent | 3,000 ÷ 1.5 = 2,000 queries a second per partition |
| Reads over the ceiling | 66,667 ÷ 2,000 | 33 times over |
| Likes over the ceiling | 10,000 small writes ÷ 1,000 | 10 times over |
The numbers come from step 2.2 of the news feed loop, which met this exact spike.
The table is 20% busy, yet reads and writes for this one post are throttled, and the throttling spreads to other posts that share its partition. Would doubling the table's capacity help? What would you change for the reads, and what for the writes?
The big picture
Synthesizing vector architecture diagram...
What to notice: the routers find each key's shard through a cached directory. One post, p9, lives on S2, so S2 alone is far over its ceiling while the other three sit at 20%. Everything on this page either spreads that load, absorbs it in the routers, or moves the piece.
What you'll be able to do after this page
- Explain how a request finds its shard: key, slot, directory, and the shard's own slot table (Part 1).
- Choose a partition key and a shard count, and tell key skew from load skew (Part 2).
- Find a hot key, and explain why adding shards doesn't help it (Part 3).
- Tame a read-hot key with a short cache, single-flight and coalescing routed by key (Part 4).
- Tame a write-hot key with salting or pre-aggregation, and compare every hot-key fix on equal terms (Part 5).
- Add capacity with
hash mod N, a ring with virtual nodes or fixed slots, and say what each moves (Part 6). - Say what DynamoDB and Kinesis split for you, and why they never split one key (Part 7).
- Move live data to a new shard without losing a write, and roll the move back (Part 8).
- Run queries and transactions that span shards, and say what they cost (Part 9).
- Trace one like and one read through every layer, and size the router fleet per AZ (Part 10).
- Map all of it to AWS services, and name the look-alikes that are not sharding (Part 11).
You may have arrived from a step that relies on this: step 2.2 of the news feed loop (the celebrity post), step 2.5 of the ad-click aggregation loop ("split the hot Kinesis shard"), step 2.1 of the digital wallet loop (1,024 buckets on 16 shards) or step 3.2 of the Discord case study (request coalescing). This page is the "why" behind all four.
Part 1. What a shard is
One machine can't hold or serve all the data, so we split it. Then every request has a new question to answer before it can do anything: which piece holds my key, and which machine holds that piece?
Key, slot, shard
A shard (DynamoDB and Kafka say partition) is one piece of the data, served by one machine or one small replica group. We don't map keys straight to shards. We map them in two steps:
- Key → slot. A fixed function turns the key into one of a fixed number of slots (also called buckets, hash slots or token ranges). Valkey and Redis Cluster use
slot = CRC16(key) mod 16,384. The slot count never changes. - Slot → shard. A directory (a small, versioned table) says which shard owns each slot. This mapping does change when we add shards or move data.
A hash tag lets several keys share a slot on purpose. If a key contains {...}, only the text between the first { and the first } after it is hashed (if there is at least one character between them). So post:{p9}, likes:{p9}:r0 and comments:{p9} all hash p9 and land in the same slot, on the same shard, where one request can read them together.
Each shard is a leader plus replicas in other Availability Zones. How those copies stay in step is the Replication, Quorums & Read-Your-Writes loop primitive's subject; on this page a shard is one unit with one ceiling.
Synthesizing vector architecture diagram...
What to notice: only the tag p9 is hashed, so every key of post p9 lands in slot 12,067. The directory names the owner; the owner's own slot table confirms it before serving.
The router, the directory and the slot table
A router is whatever computes the slot and sends the request: a proxy tier, or a "smart" client library. Routers cache the directory, because asking the directory on every request would make it the bottleneck. A cached copy can be old. So correctness can't rest on the router's copy:
- The directory is the plan: slot range → shard, plus an epoch per slot that goes up every time the slot changes owner, and a version for the whole map.
- Each shard's slot table is the truth about that shard: for every slot, "owned at epoch e", "migrating", or "moved to S4 at epoch e". A shard checks every request against it.
textROUTE(request for key k) 1. tag = text inside the first {...} of k, else all of k 2. slot = CRC16(tag) mod 16,384 3. shard = cached_directory[slot] 4. send (request, slot, epoch from the cached directory) to shard 5. the shard checks its own slot table for slot: owned (request's epoch older: also tell the router to reload) -> serve moved to S', epoch e' -> reply MOVED slot S' e' migrating -> reply "slot migrating, retry" 6. on MOVED: reload the directory, go to 3 on "retry": wait (10, 20, 40 ms ...), go to 3
That one check in step 5 is what makes moving data safe later (Part 8): a router with an old map is bounced, however old its map is.
Three words for one idea
| Word | Where you'll see it | Same idea |
|---|---|---|
| Shard / partition | Vitess, MongoDB, Citus, Kinesis / DynamoDB, Kafka | One piece of the data with its own ceiling |
| Slot / bucket / token range / vnode | Valkey, the wallet and hotel loops / Cassandra, Dynamo | The fixed unit that keys hash into and that moves as a whole |
| Router / proxy / smart client | Vitess's VTGate, Citus's coordinator, Valkey client libraries | Whatever computes the slot and sends the request |
The example we follow
A photo app keeps posts, their like counters and their comment threads in a sharded key-value store, the post store. Sixteen posts, p1 to p16, live on four shards, S0 to S3, behind six routers. At 20:01 the celebrity nova promotes her post p9. We follow that one post through every Part. The example's numbers are smaller than the hook's, so each shard's arithmetic fits on a line. Every trace on this page comes from running a private reference implementation of this store, not from working by hand.
Placement uses the real slot function. Shards own contiguous slot ranges: S0 owns 0 to 4,095, S1 4,096 to 8,191, S2 8,192 to 12,287, S3 12,288 to 16,383. From the reference implementation, the 16 posts land four per shard:
| Shard | Posts (slot) |
|---|---|
| S0 | p15 (2,138), p11 (2,270), p3 (3,689), p7 (3,821) |
| S1 | p14 (6,267), p10 (6,399), p2 (7,752), p6 (7,884) |
| S2 | p13 (10,396), p1 (11,819), p5 (11,951), p9 (12,067) |
| S3 | p16 (14,393), p12 (14,525), p4 (16,014), p8 (16,130) |
Per-like records live elsewhere. Each like also writes a record "user u liked post p" to a separate table, liked, keyed by the user (liked:{user_id}, one row per user and post), with its own shards. So a like on p9 never writes a per-event record to S2; the post store counts only post, counter and comment writes. A like is acknowledged only once it is durable in liked; the like-events stream is fed from liked's change stream (Change Streams & the Transactional Outbox), so every durable like has its event (Part 5 shows why).
| Setting | Tiny example | At real scale |
|---|---|---|
| Shards | S0 to S3, each a leader plus 2 replicas in 3 AZs. S4 is added in Part 8 as p9's own shard; the growth replay (Part 6) adds S5 on a copy | Tens to thousands of partitions |
| Shard ceiling (labelled; it mirrors a DynamoDB partition) | 1,000 writes/s and 3,000 reads/s. Reads go to the leader only, so 3,000 is the leader's; above the ceiling, requests are throttled | DynamoDB: 3,000 read units and 1,000 write units a second per partition; Kinesis: 1 MB/s or 1,000 records/s per shard |
| Slot function | CRC16(tag) mod 16,384, with hash tags | Valkey and Redis Cluster the same; DynamoDB "an internal hash function"; Kafka toPositive(murmur2(key)) % numPartitions; Kinesis MD5 of the partition key into a 128-bit range |
| Directory | Slot range → shard and epoch, versioned; v7 when Part 8 starts | Vitess keyspaces, Citus metadata, the wallet and hotel loops' bucket maps |
| Slot table on each shard | Every shard keeps its own table (slot → owned, migrating or moved-to, epoch) and checks every request | Valkey's per-node slot map; the hotel and wallet loops' "moved" marker |
| Routers | 6 (2 per AZ), stateless, each caching the directory; they learn a new version by push, or by polling every 30 s | A proxy tier or a client library |
| Router capacity (labelled) | 2,500 requests/s each | Load-test your own |
| Baseline traffic | Each post: 50 likes/s and 100 reads/s; so each shard: 200 writes/s (20%) and 400 reads/s (13%) | |
| The spike (from 20:01:00, peak by 20:02:00) | p9: 2,400 likes/s and 9,000 reads/s, replacing its baseline | The hook: 10,000 likes/s and 66,667 reads/s |
| Read service time on a healthy shard (labelled) | 5 ms | Engine- and size-dependent |
A read of p9 | One multi-key request: post:{p9} plus its counter sub-keys, all in one slot |
Snapshot S1, at 20:00:00, from the reference implementation:
| Shard | Writes/s | Reads/s | % of ceiling (w / r) | Posts | Slot table for 12,067 |
|---|---|---|---|---|---|
| S0 | 200 | 400 | 20% / 13% | 4 | — |
| S1 | 200 | 400 | 20% / 13% | 4 | — |
| S2 | 200 | 400 | 20% / 13% | 4, including p9 | owned, epoch 7 |
| S3 | 200 | 400 | 20% / 13% | 4 | — |
The story has ten beats, each in its own Part: the key is chosen (Part 2), p9 goes viral (3), the reads are tamed (4), the writes are tamed (5), the cluster grows, on a copy (6), the managed services, on a copy (7), p9 moves to its own shard while it takes writes (8), a query and a transaction cross shards (9), one like end to end (10), and on AWS (11). The full event table, including side rows left out of the body, is in Part 12.
What to remember from Part 1
- A key maps to one slot and one shard, and never gets more than that shard can serve.
- Routers cache the map; each shard owns the truth about its slots and checks every request.
- Everything later is about choosing that mapping, living with its hot spots, and changing it.
Part 2. Choosing a partition key
The key decides which queries touch one shard and which touch all of them, and whether the load spreads. Get it wrong and no amount of hardware fixes it later.
Trace: the same traffic under three layouts
Baseline traffic (800 likes a second across the 16 posts), from the reference implementation:
| Layout | Key | Where 800 likes/s go |
|---|---|---|
| Range by time | likes#<minute> | Every like of the current minute goes to the shard that owns the newest range: 800 writes/s on one shard (80%), three shards idle, and the hot spot jumps at every minute boundary |
| Hash by post | {p} | 200/s on each shard (20%) |
| Directory by creator | four creators (ada, ben, nova, omar), four posts each; nova owns p9 to p12 | Each creator is one directory row placed as a unit: 200/s on each shard today, but nova's spike will land on one shard with all her posts |
Range, hash, directory
| Layout | How keys map | Good at | Breaks when |
|---|---|---|---|
| Range | Sorted key ranges, each owned by a shard | Ordered scans ("everything from 10:00 to 10:05") | Keys grow in one direction (time, sequence IDs): every new write goes to the last range |
| Hash | A hash of the key picks the shard | Spreading keys evenly; point lookups | Any query by range or by another attribute asks every shard |
| Directory | A lookup table says where each unit lives | Placing units by decision: a big tenant on its own shard, a hotel's rooms together | You must keep the table highly available and cached |
Most real designs combine the last two: hash keys into many fixed buckets, then map buckets to shards through a directory. The wallet loop hashes wallets into 1,024 buckets on 16 shards (step 2.1), the hotel loop does the same with 1,024 buckets on 4 clusters (step 2.4), and Uber's Schemaless has 4,096 fixed shards placed on storage clusters (the Uber case study, step 2.3). Hashing spreads the keys; the directory lets us decide where each bucket goes.
Synthesizing vector architecture diagram...
What to notice: keys never rehash. Moving capacity means changing directory rows, and a single bucket (here 717, holding a big customer) can be carved out to a shard of its own.
User, tenant and time keys
- User or entity keys (
user_id,post_id,hotel_id) are the default: most requests carry one, and one entity's data stays on one shard. - Tenant keys keep a customer's data together, so per-customer queries touch one shard. The cost is that a whale customer is one unit: the drill One Customer, One Shard, One Outage is exactly this, and Part 5 and Part 8 answer it.
- Time never leads the key. A key that starts with the time sends every new write to one place. Two fixes, both in the loops: bucket a growing key by time inside an entity key, as Discord does with
(channel_id, bucket)in 10-day buckets, which bounds each partition's size (Discord case study, step 1.3); or shard a time key into sub-keys, likeHOLD#<minute>#<0-63>in the notification loop (step 2.2) and<minute>#<shard>in the job scheduler loop (steps 1.1 and 2.1), so the current minute's writes land on many shards.
A hotel booking system shards reservations by date, so "all bookings for June 3" is one shard. Why is that a bad key?
The key follows the most frequent query and the transaction boundary
| Query | Hash by post | Range by time | Directory by creator |
|---|---|---|---|
"p9's counter" | One shard | Every shard | One shard |
| "Likes in the last minute" | Every shard | One shard | Every shard |
"All of nova's posts" | Every shard | Every shard | One shard |
Anything a layout can't answer on one shard becomes a scatter-gather (Part 9). So list the top queries and the transactions first, and pick the key they already carry. The payment loop shards by merchant (step 3.1) because a merchant's ledger is what changes together.
Key skew vs load skew
- Key skew: some shards hold more keys than others. Hash by post has none here: exactly 4 posts per shard.
- Load skew: some shards take more traffic than others, because some keys are busier. Hashing does nothing about it: in Part 3, S2 holds 4 posts like everyone else and takes 12.75 times the writes of each other shard.
The loops say the same thing in their own domains: "hashing spreads symbol counts evenly, not load" in the matching engine loop (step 2.7), which assigns symbols by measured load instead; and Uber's city-sharded dispatch struggled because cities differ so much in size and spikiness (Uber case study, step 1.1).
How many shards, how many buckets
Take the larger of two numbers, and round up:
- Load: peak writes ÷ (one shard's ceiling × a target such as 0.7).
- Size: total data ÷ the most you are willing to copy or restore for one shard.
Then choose far more buckets than shards, because the bucket count caps how far you can grow without rehashing every key. Worked, from the digital wallet loop: 38,750 transfers a second at the ceiling ÷ (4,000 per shard × 0.7) = 38,750 ÷ 2,800 = 13.8, so at least 14 shards; the loop takes 16 so that 1,024 buckets divide evenly, 64 per shard.
What to remember from Part 2
- Pick the key your most frequent query and your transactions already carry.
- Hashing spreads keys, not load; a time-ordered key sends every new write to one place.
- Size shards from peak load and move size; make buckets far more numerous than shards.
Part 3. A hot key appears
At 20:02 S2 is at 255% of its write ceiling while the cluster as a whole is 79% busy. A hot key is one key whose demand is more than its shard has left. It is invisible in the average and fatal to its shard.
Trace: p9goes viral
Synthesizing vector architecture diagram...
Synthesizing vector architecture diagram...
Snapshot S2, from the reference implementation. What to notice: one bar crosses the 100% line in each chart, and it is the same shard. The other three are unchanged.
| # | Time | Event | Result |
|---|---|---|---|
| 1 | 20:00:00 | Baseline | Every shard at 200 writes/s (20%) and 400 reads/s (13%) |
| 2 | 20:01:00 | nova promotes p9 | By 20:02:00, 2,400 likes/s and 9,000 reads/s on p9 |
| 3 | 20:02:00 | S2 carries p1, p5, p13 at baseline plus p9 | Writes 150 + 2,400 = 2,550/s (255%); reads 300 + 9,000 = 9,300/s (310%). The cluster: 3,150 of 4,000 writes/s (78.75%), 10,500 of 12,000 reads/s (87.5%) |
Collateral damage. S2 throttles by shard, not by key: p1, p5 and p13 are rejected too, although their own traffic never changed. Their owners did nothing wrong except share a shard with p9.
The cluster is 79% busy on writes and S2 is at 255%. Which number matters, and would going from 4 shards to 8 help?
Why more shards don't help
| # | Replay | Result, from the reference implementation |
|---|---|---|
| 5 | Split each shard's slot range in two: 8 shards of 2,048 slots | p9 (slot 12,067) goes to the new shard 5, and so do p1 (11,819), p5 (11,951) and p13 (10,396): shard 5 takes 2,550 writes/s and 9,300 reads/s, exactly S2's load. The key can't be split, and here even its neighbours happened to follow it |
Whatever the layout, p9 alone is 2,400 writes and 9,000 reads a second on one shard: 240% and 300% of a ceiling that no fleet size changes. The same "add nodes doesn't split one key" lesson is in the key-value store loop (step 2.7), the message queue loop (step 2.7), the rate limiter loop (step 2.3), Discord (step 2.1) and Uber (step 2.1: "an airport's S2 cell is still one key on one owner").
Finding the key
Two measurements, in this order:
- Per-shard utilization against its ceiling tells you where: S2 at 255% and 310%.
- Per-key counts tell you which. Routers sample requests and keep a top-K per shard. From the reference implementation,
p9is 94.1% of S2's writes (2,400 ÷ 2,550) and 96.8% of its reads (9,000 ÷ 9,300); a 1% sample over 5 seconds already put it at 118 of 127 sampled writes.
On DynamoDB, CloudWatch Contributor Insights does step 2 for you. It has two modes: throttled keys (only the most throttled items, cheap enough to leave on) and accessed and throttled keys (also the most accessed items). One blind spot to know: write throttling caused by a global secondary index that lacks capacity isn't counted in the throttled-items graph. Watch the index's own most-accessed graph for that (Part 7).
textEVERY SAMPLED REQUEST at a router (1 in 100) 1. topk[shard][key] += 1 (a bounded top-K sketch per shard) 2. every 5 s: for each shard over 70% of a ceiling, report the keys whose share of that shard is over 10% and whether the share is reads, writes or both EVERY REQUEST at a router (admission) 3. if key has a token bucket and it is empty -> 429, Retry-After 4. if the tenant's bucket is empty -> 429, Retry-After 5. else take a token from each and route (Part 1)
Reads or writes?
Read-hot and write-hot keys need different fixes, so classify first. p9 is both: reads are 300% of the leader's ceiling (Part 4 fixes them), and writes are 240% (Part 5). A fix for one does little or nothing for the other.
Little's law says why a hot shard hurts so fast: requests in flight = arrival rate × time each spends in the system. At 9,300 reads a second and 5 ms, S2 holds about 46.5 requests at once; if throttling and queuing stretch that to 500 ms, it holds 4,650, and every client waiting on p1, p5 and p13 queues behind them.
A limit per key
The first guard is to refuse load a shard can't take, before it takes the shard down. Routers enforce budgets per key and per tenant:
| Limit | Example value (labelled) | What the caller gets |
|---|---|---|
| Per key, writes | 800/s (80% of a shard) | 429 Too Many Requests with Retry-After |
| Per key, reads | 2,400/s (80% of the leader) | 429 with Retry-After, or a cached answer if one exists |
| Per tenant | A share of the fleet set by plan | 429 for that tenant only |
| Very hot keys | Lease a block of tokens to each router so the check is local | As above; fewer calls to the limiter (the rate limiter loop, step 2.3) |
A limit protects the neighbours; it doesn't serve the celebrity's fans. Parts 4 and 5 do that. Isolating whole tenants in cells, and shuffle sharding, are the Multi-Region Failover loop primitive's subject (coming).
What to remember from Part 3
- A hot key is one key's demand above its shard's ceiling; the cluster average hides it.
- Adding shards moves other keys, never this one.
- Find it per key, then ask: hot for reads, writes, or both?
Part 4. Hot reads: cache, coalesce, copy
p9 gets 9,000 reads a second, and they are all the same question: "what does post p9 look like now?" Asking the shard 9,000 times a second for one answer is the waste to remove.
Cache it for a second
Each router keeps a micro-cache: p9's answer for 1 second, with two extras from the news-feed loop's step 2.2:
- Single-flight: on a miss, only one fetch goes to the shard; every other request for the key waits for that fetch.
- Stale-while-revalidate: an entry older than 1 second but younger than 3 seconds is still served, while one refresh runs in the background.
textREAD(k) at a router 1. entry fresh (age < 1 s) -> serve it 2. entry stale but age < 3 s -> serve it; if no refresh is running, start one 3. no usable entry: a fetch for k is in flight -> wait for it else -> start the fetch; everyone waiting gets its result 4. a write to k seen by this router -> drop the entry
| # | Time | Event | Result, from the reference implementation |
|---|---|---|---|
| 6 | 20:02:10 | Micro-cache on | Each router fetches p9 once a second: S2 sees 6 requests/s for p9 instead of 9,000 (each one multi-key request for the slot) |
| 7 | 20:02:15.000 | nova edits p9's caption on S2 | Each router shows the old caption until its next refresh: the first reader to see the new caption came between 0.187 s (R2) and 0.779 s (R1) later. The bound is 1 s plus one query; if refreshes fail, stale-while-revalidate can serve the old caption for up to 3 s |
The staleness is now a stated number, not an accident. A cache in front of a hot key must have one: a short TTL, or an invalidation path that every write takes. The news feed loop's step 2.2 met the counter-example: the DAX query cache isn't invalidated by writes to the table.
Why not cache it forever?
Because then a caption edit, a takedown or a price change is invisible until the entry is evicted: a write you can't see for a day. And a cache restart is the worst moment: every entry is gone at once, so thousands of readers miss together and each would query the shard. Single-flight is what saves that moment, the cold-cache stampede in the drill The Product Page That Melted Redis: each router sends one fetch per key and everyone else waits for it. A short TTL with single-flight bounds both the staleness and the stampede.
Coalesce identical reads
Some reads can't be cached, even for a second (a check that must see the latest write). Request coalescing still helps: while a query for p9 is in flight, identical requests join it instead of starting their own. A request may join only if no write to the key completed since the query began, or it could get an answer older than a write it already saw acknowledged.
With requests arriving at rate λ and each query taking t, the shard sees:
The formula is exact for random (Poisson) arrivals and a fixed query time; the reference implementation's simulation matched it to within 0.2% (176.25 against 176.47, 195.61 against 195.65).
| # | Setup | Arithmetic | Reads/s on S2 |
|---|---|---|---|
| 6a | Coalescing on each of 6 routers, requests spread at random: 1,500/s each, t = 5 ms | 1,500 ÷ (1 + 1,500 × 0.005) = 1,500 ÷ 8.5 = 176.5 per router, × 6 | 1,059 |
| 6b | Routed by key to one coalescer | 9,000 ÷ (1 + 9,000 × 0.005) = 9,000 ÷ 46 | 196 |
| 6b | Routed by key to one coalescer per AZ (3 rings, 3,000/s each) | 3,000 ÷ (1 + 15) = 187.5, × 3 | 562.5 |
Synthesizing vector architecture diagram...
What to notice: three requests, one query. The next request after a completed write can't join an older query, so it starts a new one.
Coalescing on each of 6 routers cut 9,000 reads to 1,059. Why did routing by key cut them to 196?
Coalescing isn't a cache. A coalesced result lives only while its query is in flight, 5 ms here. It never serves an answer older than one query, and it does nothing for requests that don't overlap. A micro-cache serves one answer for a whole second. Use the cache where a second of staleness is fine, and coalescing where it isn't.
The edge can do the same in front of the routers: a CDN that collapses identical requests for one object into one fetch to the origin (Part 11).
Replica reads
Letting the shard's two replicas serve reads too raises S2's read ceiling to 3 × 3,000 = 9,000 a second. S2 needs 9,300: still 103% (event 7a in Part 12). And replica reads are stale by the replica's lag (page 06); they raise the ceiling of the whole shard by 3 times, whatever the key; and they do nothing for writes, which every copy must still apply. Useful for a warm shard, not a cure for one hot key.
Copies of one value
| Fix | How | Reads | Writes | Staleness |
|---|---|---|---|---|
| Read copies (event 7b) | 8 copies under post:{p9#c0} to post:{p9#c7}, which land two per shard; a read picks one at random | 1,125 per copy; two copies per shard, so 2,250 + other reads: 2,650 (88%) on S0, S1, S3 and 2,550 (85%) on S2 | Every caption edit writes all 8 | Copies briefly differ while an edit is being written |
The search autocomplete loop uses the same idea (step 3.2, more replicas for hot shards). Every fix for hot reads is compared, with the write fixes, in one table at the end of Part 5.
What to remember from Part 4
- Cache hot reads for a short, stated time, and protect the miss with single-flight.
- Coalescing needs identical requests to meet: route by key.
- Replicas raise a whole shard's reads; copies of one value trade one write for N.
Part 5. Hot writes, and every fix on equal terms
p9's counter takes 2,400 likes a second, and one shard takes 1,000 writes. Caching doesn't help a write: each like must change the count. Either the writes go to several places, or fewer writes are sent.
Salting
Salting (write sharding) splits one key into S sub-keys and sends each write to one of them. A read adds them back up.
The naive sizing is S = rate ÷ ceiling = 2,400 ÷ 1,000 = 2.4, so 3. It forgets that a shard carries every key hashed to it. Three sub-keys of 800 a second each, on shards that already carry 200, put those shards at exactly 1,000: 100%. The right rule, in words (Formula 2): each shard's load is its other traffic plus, for each sub-key it holds, the key's rate ÷ S; choose S and the sub-keys' placement so every shard stays at or under 80%. Here each non-S2 shard has 800 − 200 = 600 a second to spare, so the 2,400 must be spread evenly over all four shards.
We take S = 8, with the salt inside the hash tag: likes:{p9#0} to likes:{p9#7}. Each sub-key takes 2,400 ÷ 8 = 300 likes a second.
| # | Time | Event | Result, from the reference implementation |
|---|---|---|---|
| 8 | 20:02:20 | Salting, S = 8 | The sub-keys land on slots 15,527, 11,398, 7,397, 3,268, 15,395, 11,266, 7,265 and 3,136: two per shard. S2: 150 + 2 × 300 = 750 (75%); S0, S1, S3: 200 + 600 = 800 (80%) |
Synthesizing vector architecture diagram...
What to notice: every shard carries two sub-keys, and each shard's total counts its other traffic, not just its share of p9. (Arrows to the second sub-key of each shard are left out to keep the picture readable.)
On four shards, S = 4 would do the same: suffixes #0 to #3 already land one per shard, and give the same 750 and 800 with half the read fan-out. Loops pick S larger (10 to 64) to leave room for growth and for collisions they can't see. Changing S later means every reader must learn the new S first, and for an index whose reader looks up one member (the news feed loop's FOLLOWED_BY#<id>#<n>, step 2.1), a backfill too.
When salts collide, and the tag trap
| # | Replay | Result, from the reference implementation |
|---|---|---|
| 8a | The salts land at random instead of where we checked | Of the 4^8 = 65,536 ways 8 salts can fall on 4 shards, only 2,520 put exactly two on each. So 1 − 2,520 ÷ 65,536 = 96.15% of the time some shard gets 3 or more (a Monte Carlo run of 200,000 trials gave 96.12%). A shard with 3 salts: 200 + 3 × 300 = 1,100 (110%) |
| 8a | The salt outside the tag: likes:{p9}#3 | Only p9 is hashed, so all 8 sub-keys land on slot 12,067, on S2: nothing was spread. Put the salt inside the tag: likes:{p9#3} |
So: where you know the slot function (Valkey, your own directory), check where each sub-key lands and pick suffixes that spread. Where you can't see partitions (DynamoDB), leave headroom for collisions with a larger S. DynamoDB's documentation also describes calculated suffixes, which solve a different problem: the suffix is derived from an attribute of the item, so a reader of one item knows which sub-key to read, without reading all S.
Salting cut S2's writes to 75%. Why is every shard now at 600% of its read ceiling, and what fixes it?
| # | Replay | Reads reaching the shards |
|---|---|---|
| 8b | Count read uncached, S = 8 | 72,000 sub-key reads/s: 18,000 per shard (600%) |
| 8b | Count read through the 1 s micro-cache | 6 routers × 8 = 48 sub-key reads/s, plus 6 reads of post:{p9} |
Pre-aggregation
Most of those 2,400 writes are "+1". Each router can add them up and write the total. That is pre-aggregation (write sharding with local sums):
| # | Time | Event | Result, from the reference implementation |
|---|---|---|---|
| 9 | 20:03:00 | Switch the counter to pre-aggregation | Each router takes about 2,400 ÷ 6 = 400 of p9's likes a second, sums them for 500 ms and writes its own sub-key likes:{p9}:rN: 6 routers × 2 flushes = 12 writes/s on slot 12,067. S2: 150 + 12 = 162 writes/s |
(A store with conditional writes can instead take one conditional write per window: each router adds its 500 ms sum under "apply only if this window's id is new", so a retry can't count twice. Our per-router sub-keys do the same with seq.) The switch is itself a small migration: the eight salted sub-keys stop taking writes, and a read of the count adds their frozen sum, which never changes again, to the six router sub-keys.
textPRE-AGGREGATE p9's likes, on router N 0. (precondition) nothing else writes the hot key once per event on each like: 1. write "u liked p9" to liked:{u} (durable); like-events is fed from liked's change stream 2. acknowledge the like 3. buffer += 1 every 500 ms: 4. total = total + buffer; seq = seq + 1; buffer = 0 5. write likes:{p9}:rN = (total, seq), applied only if seq > the stored seq on restart: 6. read likes:{p9}:rN back: total and seq continue from there a read of the count: 7. sum likes:{p9}:r0 .. r5 (one multi-key request, one slot) plus the frozen salted sum
Nothing else may write the hot key
Pre-aggregating a counter only helps if nothing else writes the hot key per event. If each like also wrote a record "u liked p9" under p9's tag, S2 would still take 2,400 writes a second and the counter's 12 would change nothing. That is why the per-like records are keyed by the user, liked:{u}, in a separate table: they spread over that table's shards by user, and "who liked p9?" becomes a derived index or a scatter (Part 9). The same shape shows up in the mobile news feed loop (R2.5), which sums counters from the change stream instead of incrementing per like.
When a router crashes
A flush can fail, be retried, or never happen. The design above makes each case harmless:
| # | Time | Event | Result, from the reference implementation |
|---|---|---|---|
| — | 20:03:05.210 | R2 flushes likes:{p9}:r2 = 2,059, seq 11 | Applied |
| — | 20:03:05.410 | R2's client timed out and retries the same flush | Ignored: seq 11 is not larger than the stored 11 |
| 10 | 20:03:07.300 | R4 crashes just before a flush, with 211 likes counted but not flushed; its last flush (20:03:06.800) wrote 2,772, seq 14 | Its sub-key still says 2,772. Its clients' next likes go to the other routers |
| — | 20:03:09.000 | R4 restarts and reads back 2,772, seq 14 | Its next flush, at 20:03:09.300, writes 2,862, seq 15 |
The 211 likes are not lost: each was durable in liked:{u} before its "OK", and reaches like-events through the table's change stream. The count is 211 low (the reference implementation checked: short by exactly R4's lost buffer) until the job that rebuilds each post's count from like-events every few minutes (our design choice) corrects it: it recounts p9's likes and writes the difference from the sub-keys' sum into a correction sub-key likes:{p9}:fix, which reads also add. The count is a view, rebuildable from durable events. Why the flush must carry an absolute total and not "+211" is the Idempotency & Effectively-Once Processing loop primitive's rule (Part 5): an additive flush that is retried counts twice.
When the pre-aggregation runs inside a stream job instead of a router, it has more rules (lateness decided per event, partials flushed before a watermark and at a checkpoint barrier); those belong to the Event Time, Watermarks & Checkpoints loop primitive (Part 3). The ad-click loop (step 2.5) and the YouTube loop (step 2.5) both pre-aggregate this way.
Leasing and striping
- Token leasing is pre-aggregation for a limiter: instead of decrementing a hot tenant's counter on every request, each server leases a block of tokens and spends them locally (rate limiter loop, step 2.3).
- Striping a balance handles a hot account that has a floor (it may not go below zero), where salting's "read all S and add" would let two stripes spend the same money. Credits land on stripes; a sweeper moves stripes into the main account; debits come only from the main account, where the floor is checked under one lock. S is sized from the lock hold time (digital wallet loop, step 2.3; payment loop, step 3.1). The cost is delay: a credit is spendable only after the sweep.
Counters that must merge concurrent increments without coordination can also be CRDT counters, one sub-count per writer; that is the same idea as the per-router sub-keys.
A dedicated shard
Moving p9's slot onto a shard of its own protects p1, p5 and p13, but before the fixes that shard alone would carry 2,400 writes (240%) and 9,000 reads (300%): isolation, not a cure. A dedicated shard can't give one key more than one shard's ceiling. With the fixes on, it is worth doing, and Part 8 does it, live. For a customer you know is coming (the drill's 10× tenant), carve out its buckets before launch day.
GSI write sharding
A hot key in a secondary index is as dangerous as one in the table. In DynamoDB, a global secondary index that can't keep up throttles writes to the base table. If an index is keyed by creator_id, every write about any of nova's posts lands on one index key: Part 7 shows 2,550 a second. The fix is the same salting, on the index key: creator_id#0 to creator_id#7, with reads querying all 8 in parallel and merging (URL shortener loop, step 2.5: owner_id#0..15).
Every hot-key fix on equal terms
Each row is filled from both sides: does a read fix help writes, and does a write fix help reads? Numbers are p9's, from the reference implementation.
| Fix | Reads reaching the shard | Writes reaching the shard | Staleness | Added write delay | Read amplification | On a crash | Values it works for | What it needs | Operational cost |
|---|---|---|---|---|---|---|---|---|---|
| Micro-cache + single-flight | 6/s (one per router per second) | Unchanged: 2,400/s | Up to 1 s (3 s if refreshes fail) | None | None | Cache refills; single-flight stops the stampede | Anything readable a second late | A TTL or an invalidation path | Router memory |
| Coalescing, routed by key | 196/s (562.5 per AZ) | Unchanged | At most one query (5 ms) | None | None | In-flight queries retried | Anything, even must-be-fresh reads | Routing by key; a per-key write counter | A hop; ring changes briefly split a key |
| Replica reads | Spread over 3 copies: ceiling 9,000/s, still 103% | Unchanged; every copy applies every write | Replica lag | None | None | A replica's failure removes a third of the read ceiling | Reads that tolerate lag | Replicas; lag monitoring | Raises the whole shard, not the key |
| Read copies (8) | 1,125/s per copy; 2,250 per shard plus its other reads (85 to 88%) | Every edit × 8 | Copies briefly differ | The write of all 8 | None | A half-written edit leaves copies different | Read-mostly values | Fan-out on write | 8 keys per hot value |
| Salting (S = 8) | × 8 per read (72,000/s) unless the sum is cached (48/s) | 300/s per sub-key; 750 to 800 per shard | None | None | × S | A lost write is lost; a retried +1 counts twice unless idempotent | Mergeable values only | Checked placement or collision headroom | Readers must know S |
| Pre-aggregation | Unchanged (pair it with a cache) | 12/s | Count lags up to 500 ms | Up to 500 ms | × 6 sub-keys, one request | Unflushed sums re-derived from durable events; retries ignored by seq | Mergeable values only | Nothing else writing the key per event; durable events | A rebuild job; one sub-key per writer |
| Dedicated shard | Up to 3,000/s (the whole shard) | Up to 1,000/s (the whole shard) | None | None | None | As any shard | Anything | A live move (Part 8) | A shard per whale |
| Per-key limit | Capped at 2,400/s; the rest refused | Capped at 800/s; the rest refused | None | None | None | Refused requests retried by callers | Anything | Callers that handle 429 | Unhappy fans |
Read the table down the "writes" column: nothing for reads touches writes. Read it down the "reads" column: salting makes reads worse unless you cache the sum. That is why p9 needs one fix from each half.
Snapshot S3, after event 9 (micro-cache and pre-aggregation on), from the reference implementation:
| Shard | Writes/s | Reads/s (requests) | % of ceiling (w / r) |
|---|---|---|---|
| S0, S1, S3 | 200 | 400 | 20% / 13% |
| S2 | 150 + 12 = 162 | 300 + 6 = 306 | 16% / 10% |
What to remember from Part 5
- Salting splits the writes and multiplies the reads: size it after the shard's other traffic, with the salt inside the tag.
- Pre-aggregation writes totals, not events: nothing else may write the hot key per event, and the flush must be safe to repeat.
- A dedicated shard protects the neighbours; it can't make one key bigger than a shard.
Part 6. Adding capacity: rings, slots and splits
Day 2, on a copy of the cluster: the app is growing, and every post's baseline has doubled to 100 likes and 200 reads a second, so each shard is at 400 writes a second (40%). We add a fifth shard, S5. The question is how much data moves, and whether we get to choose which.
Trace: three ways to add S5
From the reference implementation, over 10,000 synthetic keys and our 16 posts:
| Method | How it places a key | Keys that move, 4 → 5 shards | Who chooses |
|---|---|---|---|
hash mod N | shard = hash(key) mod N | A key stays only if its hash mod 4 equals its hash mod 5: 4 of every 20 values. 80% move (the run: 80.3% of the keys, 13 of the 16 posts) | Nobody: the arithmetic decides |
| Consistent hashing, 16 random tokens per shard | Each shard puts 16 points (tokens) on a ring; a key goes to the next token in ring order | 17.2% in our run (the median over 30 random layouts: 19.9%) | The random tokens |
| Fixed slots + directory | 16,384 slots never change; S5 takes chosen slots | 16,384 ÷ 5 = 3,276.8, so S5 takes 3,277 slots, 819 or 820 from each old shard: 19.8% of the keys | We do |
hash mod N
mod N is the naive answer every loop warns against: the S3-like storage loop (step 1.1: hash % 100 → 101 moves almost everything) and Uber's dispatch (step 2.1: adding the 301st machine moves 300 of every 301 keys).
| # | Change | Keys that stay | Moved |
|---|---|---|---|
| K7 | 4 → 5 | h mod 4 = h mod 5: 4 of 20 residues | 80% |
| K7b | 4 → 8 (doubling) | h mod 4 = h mod 8: 4 of 8 residues | 50% (the run: 50.3%) |
| K7b | 12 → 16 | h mod 12 = h mod 16: 12 of 48 residues | 75% (the run: 75.2%) |
Doubling is the best mod N can do, which is why "double the partitions" is a common step. It still moves half the data, and for an ordered topic it moves half the keys' new records to a different partition (Part 7).
The ring, with one point and with many
With one token per shard, each shard owns one arc of the ring, and random arcs are very uneven. Adding S5 cuts one arc in two:
Synthesizing vector architecture diagram...
Arcs in ring order, positions as fractions of the ring, from the reference implementation. A key belongs to the next token at or after its hash. What to notice: S5's token lands inside S3's arc, so S5 takes everything between 0.037 and 0.381 from S3, and S0, S1 and S2 shed nothing.
| # | Ring | Before (keys per shard) | After adding S5 | Moved |
|---|---|---|---|---|
| K7a | 1 token per shard | S0 2,965, S1 2,281, S2 924, S3 3,830: the busiest holds 1.53 times the average | S3 keeps 383; S5 takes 3,447, all from S3; the busiest is 1.72 times the average | 34.5% |
| K7 | 16 tokens per shard | S0 3,046, S1 2,809, S2 1,793, S3 2,352: 1.22 times the average | S5 takes 1,715: 1,388 from S0, 104 from S1, 195 from S2, 28 from S3; the busiest is 1.35 times the average | 17.2% |
Synthesizing vector architecture diagram...
What to notice: with many tokens, S5's points cut many different arcs short, so it takes a slice from several shards instead of all from one. With only 16 random tokens the slices are still uneven: in our run S0 gave 1,388 keys and S3 gave 28.
How many tokens is enough? Over 30 random layouts in the reference implementation, the busiest shard after adding S5 held a median of 1.26 times the average with 16 tokens, 1.11 with 100 (Uber's Ringpop default, step 2.1) and 1.07 with 256. Many small pieces are what make a move even. Cassandra defaults to 16 tokens per node, and pairs them with a token allocator that places new tokens to balance load for a given replication factor, instead of picking them at random. This answers the drill The Ring Rebalance That Crushed Node 07: without virtual nodes, a new node takes keys only from the one node after it on the ring, so the old nodes don't each shed a fair share.
Why not hash mod 5, and why not a ring with one point per node?
Fixed slots and a directory
Nothing is rehashed: keys stay in their slots, and we decide which slots S5 takes. Because we know where the load is, we choose one block of 819 or 820 slots from each old shard, each holding one of our posts:
| From | Slots S5 takes | Post that moves |
|---|---|---|
| S0 | 1,320 to 2,138 | p15 |
| S1 | 5,449 to 6,267 | p14 |
| S2 | 9,578 to 10,396 | p13 |
| S3 | 13,574 to 14,393 (820 slots) | p16 |
Snapshot S4 (day 2, on the copy), from the reference implementation: S5 takes 3,277 slots and 19.8% of the synthetic keys; with 4 posts at 100 likes a second, S5 is at 400 writes/s (40%), and S0 to S3 drop to 3 posts each, 300 writes/s (30%).
The Dynamo paper (2007) found that fixed, equal-sized partitions assigned to nodes balanced load best of the strategies it tried, and turned a move into copying whole partition files. Uber's Schemaless grows the same way: a fixed 4,096 shards, and "we typically expand by splitting each MySQL server in two" (case study, step 2.3). The wallet and hotel loops' 1,024 buckets are the same design.
Split and merge
Stores that own ranges (Bigtable-style, and DynamoDB internally) grow differently: a range that grows too big or too hot is split in two, and cold neighbours are merged. Split at the key that halves the load, not the key count. On our copy, S2's load sits in one slot, so a load split cuts at 12,067: p13, p1 and p5 (slots 10,396 to 11,951) on one side, slot 12,067 alone on the other, and no split goes finer than one slot, or one key (event K7c in Part 12). The cost is brief throttling of the range while it splits (DynamoDB: splits "usually complete in the order of minutes"). A range whose writes all go to its tail (a monotonic key) can't be helped by splitting either: the new tail takes every write. The S3-like storage loop (step 2.6) spreads such a tail over hash-prefixed sub-ranges instead.
Bounded load, briefly
Consistent hashing with bounded loads caps each node at, say, 1.25 times the average load and sends overflow keys to the next node in ring order. It evens out skewed traffic for stateless placement (which server handles a request) and for owners a controller can move (the web crawler loop, step 3.1). It is the wrong tool for a store: the key's owner now depends on current load, so a spilled key's request goes to a node that doesn't have its data, and ownership moves whenever loads shift. And like vnodes, it never splits one hot key. This is the drill's second question, and Uber's answer in step 2.1.
On equal terms
hash mod N | Ring with vnodes | Fixed slots + directory | Split and merge | |
|---|---|---|---|---|
| Keys moved on adding one of N + 1 | Most (80% for 4 → 5) | About 1/(N + 1), from random places | Exactly what we choose (about 1/(N + 1)) | Only the split range |
| Who chooses what moves | Arithmetic | Random tokens (or an allocator) | The operator or a balancer, by load | The store, by size or load |
| Balance (busiest ÷ average) | Even keys | 1.26 with 16 random tokens, 1.11 with 100 (median, our runs) | As even as we place it | By load, if split by load |
| Lookup | A division | Binary search over tokens | One array index | Search over ranges |
| Metadata | None | Tokens × nodes | One row per slot | One row per range |
| One hot key | Stays on one shard | Stays on one node | Stays in one slot | No split below one key |
| Machines of different sizes | Not possible | More tokens for bigger nodes | More slots for bigger nodes | By load |
What to remember from Part 6
hash mod Nmoves almost everything; a ring moves about 1/N; fixed slots move exactly what you choose.- Many small pieces (vnodes or slots) are what make a move even.
- No split or ring goes finer than one key.
Part 7. DynamoDB and Kinesis: managed splitting never splits one key
Managed services split partitions for you. So does the celebrity post just work on DynamoDB and Kinesis? The same p9, replayed on copies of each.
DynamoDB: adaptive capacity and split for consumption
| Fact | What AWS states |
|---|---|
| The unit | A partition. DynamoDB uses "an internal hash function" of the partition key to pick it; items with the same partition key are stored together, sorted by sort key |
| Ceiling | Throttling "if a single partition receives more than 3000 read operation or more than 1000 write operations" a second |
| Burst | Unused capacity is retained for "up to five minutes (300 seconds)"; AWS notes that burst details "might change in the future" |
| Adaptive capacity | "automatically and instantly" raises a hot partition's share, as long as traffic stays within the table's capacity and the partition maximum |
| Isolating hot items | Frequently accessed items are moved apart; for "consistently high traffic to a single item", a partition may contain "only that single, frequently accessed item", served "up to the partition maximum of 3,000 RCUs and 1,000 WCUs" |
| Splitting an item collection | By sort key, unless the traffic follows a monotonically increasing or decreasing sort key; never across partitions when the table has a local secondary index (and only then does the 10 GB item-collection limit apply) |
| Split for consumption | The DynamoDB paper (USENIX ATC 2022) calls splitting a partition by its traffic a split for consumption; it takes minutes, and it doesn't help a single hot item or a sequential range of keys |
| On-demand | Instantly handles up to double the previous peak; exceeding double within 30 minutes can throttle. New on-demand tables start able to serve 4,000 writes and 12,000 reads a second |
| GSIs | Partitioned the same way; a GSI that can't keep up throttles writes to the base table |
| # | Replay | Result, from the reference implementation |
|---|---|---|
| D1 | p9's counter as one DynamoDB item, 2,400 small writes a second | Adaptive capacity isolates the item on its own partition, up to 1,000 WCU: 1,400 writes a second are still throttled (2,400 − 1,000). No split can help: the item is one key |
| D1 | A labelled side replay: a table post_likes with PK = p9, SK = user_id | Sort keys are random, so the item collection can be split by sort key and spread. With SK = liked_at, every write goes to the newest end: a monotonic sort key, which can't be split around |
D2: a hot GSI key. Suppose post_likes has a GSI keyed by creator_id. Every like of any of nova's four posts writes the index key nova: 2,400 + 3 × 50 = 2,550 writes a second on one index partition key, which throttles the index and, through it, the base table. Salt the index key: nova#0 to nova#7, 318.75 a second each, and query all 8 in parallel to read it. Contributor Insights won't show the index's write throttling in its throttled graph (Part 3); the index's most-accessed graph will.
Kinesis: shards, splits and lineage
A Kinesis stream is a set of shards. Each record's partition key is hashed with MD5 into a 128-bit integer, and the shard whose hash-key range contains that integer takes the record. Each shard takes 1 MB/s or 1,000 records/s of writes, whichever comes first, and serves 2 MB/s and 5 GetRecords calls a second of reads. The limit is shared by every key hashed to the shard.
Kinesis rejects p9's like events. Why won't splitting its shard help, and what key would you use?
| # | Replay (stream like-events, 4 shards, key post_id, ~200-byte records) | Result, from the reference implementation |
|---|---|---|
| K4 | MD5 placement of the 16 posts | 3, 4, 3 and 6 posts on shards 0 to 3: key skew from a real hash with few keys. p9 is on shard 3 with p1, p5, p6, p12 and p13 |
| K4 | The spike | Shard 3: 2,400 + 5 × 50 = 2,650 records/s against 1,000: 1,650 a second rejected. Bytes: 2,650 × 200 B = 530 KB/s, under 1 MB/s: records bind first |
| K4 | "Split that shard" at its midpoint | p9 goes to the lower child with p6 and p12: 2,500 records/s, still 2.5 times the limit |
| K4 | A random key | 3,150 ÷ 4 = 787.5 records/s per shard (79%); 157.5 KB/s |
Synthesizing vector architecture diagram...
What to notice: after a split the parent stops taking records but keeps the ones it has. A reader that starts the child holding p9 before finishing the parent can see p9's newer events before its older ones.
K6: resharding the whole stream. UpdateShardCount 4 → 8 with uniform scaling: p9 ends on shard 6 of 8. The 4 parents become CLOSED and later EXPIRED; the 8 children take new records. By default UpdateShardCount can't go above double or below half the current count in one call, or run more than 10 times per rolling 24 hours per stream. Reading children only after their parents is the Change Streams & the Transactional Outbox loop primitive's subject (Part 9).
K5: packing. Putting 9 likes in one record cuts 3,150 records a second to 350, but a record has one partition key for everything inside it, so the key can no longer be post_id. Packing only works once the key doesn't have to name an entity.
Kafka, in one sentence. Kafka picks a partition with toPositive(murmur2(key)) % numPartitions, so adding partitions sends about half the keys' new records to a new partition when doubling (the reference implementation: 50.4% of 10,000 keys going 4 → 8; p9 happened to stay in partition 2), existing records stay where they are, and the partition count can't be reduced; what that does to order is the Change Streams page's Part 4.
What to remember from Part 7
- Managed services split partitions for you, but never below one key or one item.
- A stream shard takes 1,000 records or 1 MB a second, whichever comes first, shared by every key on it.
- Adding shards or partitions sends a key's new records somewhere new: readers follow lineage.
Part 8. Moving live data
The fixes hold tonight, but nova posts daily and every spike still lands on p1, p5 and p13's shard. We give p9's slot a shard of its own, and we do it while it's being liked. That means copying data that keeps changing, switching the owner without losing a write, and coping with routers that haven't heard.
At 20:10 slot 12,067 holds post:{p9}, the six counter sub-keys and p9's comment thread: 200 MB (labelled). With the fixes on, the slot takes about 12 writes a second (the counter flushes), plus a caption edit and comments now and then; S2 as a whole takes about 162.
Copy and tail
| # | Time | Event | Result, from the reference implementation |
|---|---|---|---|
| 14 | 20:10:00.000 | Directory v7: slot 12,067 → S2, epoch 7. S4 is created. The copy starts from a replica of S2 that has applied S2's log through position 48,200 | A copy source must be caught up to a known position (page 06, Part 8) |
| 15 | 20:10:00 to 20:10:10.000 | 200 MB at 20 MB/s = 10 s, throttled so S2's replica keeps serving | The copy shares disks and links with live traffic |
| 16 | 20:10:10.000 to 20:10:10.421 | S4 tails S2's log from position 48,201, applying only slot 12,067's entries | The backlog was 1,619 log entries (121 of them for the slot); scanning at 4,000 entries/s (labelled) while new ones arrive, S4 reaches the head at 20:10:10.421. From then on each write reaches S4 about 3 ms after S2 (labelled: 1 ms network, 2 ms apply) |
| 17 | 20:10:10.421 to 20:11:10.400 | Shadow reads: every cache-miss fetch of p9 (about 6 a second) also goes to S4, and the answers are compared by log position; S2's answer is the one returned | 358 compares. In 16 of them S4 was a few milliseconds behind on a slot key; for example at 20:10:38.372 S2 answered at position 54,437 and S4 at 54,433, one counter flush behind. Rechecked once S4 reached S2's position, all 16 matched: no real mismatch |
Starting from a snapshot at a known position and then applying the log from the next position is the same bootstrap rule as the Change Streams & the Transactional Outbox loop primitive's Part 7. Instead of shadow reads, a mover can compare by scanning the whole slot on both sides at one position (Vitess calls its tool VDiff).
Synthesizing vector architecture diagram...
What to notice: S2 stops accepting the slot at the fence, before anyone owns it elsewhere, and S4 starts only after the directory commit. R5 never heard about v8, and S2's own slot table sent it to S4.
Fence and flip
textMOVE slot s from source to target (directory at version v, slot epoch e) 1. copy a snapshot of s at log position p, from a copy applied through p 2. target tails the source's log from p + 1, applying s's entries 3. shadow-read (or scan-compare) until source and target agree by position 4. FENCE: source appends "s: migrating, epoch e + 1" to its log as entry f; from now until step 5 commits, it answers every read AND write for s with "slot migrating, retry"; routers retry with backoff 5. when the target has applied through f (the fence entry itself): commit directory version v + 1: s -> target, epoch e + 1 (the commit point) target slot table: "s: owned, epoch e + 1" source slot table: "s: moved to target, epoch e + 1" (kept forever) 6. routers learn v + 1 by push or poll; any request that reaches the source gets MOVED, reloads the directory and retries (it must be safe to repeat) 7. after the compare is clean, delete the source's copy of s; the source still answers MOVED for s from its slot table 8. abort before step 5 = drop the target; after step 5, roll back by the same move in reverse 9. invariant: exactly one shard accepts reads and writes for s, and it holds every write any shard accepted for s before
| # | Time | Event | Result, from the reference implementation |
|---|---|---|---|
| 18 | 20:11:10.400 | Fence: S2's slot table says "migrating, epoch 8"; the fence is S2's log entry 59,648 (its log reached 59,647 at the fence) | Reads and writes for the slot get "slot migrating, retry". Refusing reads too means a router can't be served a frozen copy |
| 19 | 20:11:10.403 to .413 | S4 applies through 59,648 at .403; the mover, checking S4's position every 10 ms (labelled), sees it at .410; directory v8 is committed at 20:11:10.412; the two slot tables are set at .413 | The slot accepted no writes for 13 ms. No flush or refresh happened to fall in those 13 ms; a probe flush sent at .405 got "retry", tried again at .415, got MOVED (its router hadn't heard of v8 yet), reloaded and succeeded at S4 at .419: 14 ms late, well inside a 500 ms flush buffer |
| 20 | 20:11:10.418 | Five routers receive the v8 push | They send p9's requests to S4 |
The fence is the Leases, Fencing Tokens & Distributed Locks loop primitive's rule (Part 4) applied to a slot: the epoch is the fencing token, and the check happens at the place that holds the data.
The directory says S4. Router R5 missed the push and still has v7; it polls only every 30 seconds. What stops its next write from vanishing?
Synthesizing vector architecture diagram...
Snapshot S5, at 20:11:10.500. What to notice: the router with the old map reaches the old owner, and the old owner's own table turns it away; nothing depends on the router having heard about v8.
Without the slot table
| # | Counterfactual | Result, from the reference implementation |
|---|---|---|
| 21x | S2 has no slot-table check and accepts R5's flush at 20:11:10.500 into its old copy | S4 never sees it, and readers on v8 read S4. R5 keeps writing S2 until its next poll at 20:11:31.200: 42 flushes land on the old copy. The counter survives only because each flush carries the absolute total, so R5's first flush to S4 catches up; a comment posted through R5 at 20:11:12.100 lands only on S2's copy and is lost for good. A read through R5 is served S2's frozen copy |
In the hotel reservation loop's terms, that lost write is a hold taken on the old copy: an oversold room. In the digital wallet loop, it is a double spend. Both loops use this page's rule: the old shard keeps a permanent "moved" row for the bucket, checked under the same lock by every transaction, answered with a retryable error that makes the service reload the map (hotel step 2.4, wallet step 2.1 and R2.11). Uber met the same problem with ring owners during membership churn and fenced at the store with a version check (case study, step 2.7).
| # | Time | Event | Result |
|---|---|---|---|
| 22 | 20:12:10.400 | The shadow compare is clean; S2 deletes its copy of slot 12,067 | S2 still answers MOVED 12067 S4 epoch 8 for the slot, now and later |
When the move fails
| # | Branch | What happens, from the reference implementation |
|---|---|---|
| 19a | S4 dies at 20:11:10.405, after S2's fence but before v8 commits | S4 never confirms the fence; the mover gives up 100 ms after the fence (labelled), at 20:11:10.500. It commits directory v9: slot 12,067 → S2, epoch 9, at .502, and sets S2's slot table to "owned, epoch 9" at .503. S2 resumes. Epoch 8 is skipped because it was already handed out at the fence; anything still carrying epoch 8 is now refused. The only cost is a 103 ms pause |
| 17a | S4's leader dies at 20:10:05.000, halfway through the copy (Part 12) | Directory still v7, S2 never fenced: nothing to undo. The partial copy is dropped and the move restarts from a fresh snapshot |
Abort before the flip = drop the target. Rollback after the flip = the same move in reverse. Vitess keeps a reverse replication stream from the new owner back to the old one after it switches writes, so that ReverseTraffic can roll back quickly; its documentation warns that turning reverse replication off means "you will not be able to rollback once you have switched write traffic".
Dual writes instead of a log tail
Some teams keep both copies current by having the application write to both during the move, and backfill the old data separately. It works only with versions (event 23 in Part 12, from the reference implementation):
- The backfill overwrites a newer write. It read the caption at v4, a live dual write then put v5 on both shards, and the backfill's late write of v4 landed on S4: S2 says v5, S4 says v4.
- Two writers race. Routers R1 and R2 write caption v5 and v6 to both shards; S2 receives v5 then v6, S4 receives v6 then v5. Without versions, S2 ends at v6 and S4 at v5.
The cutover order: dual-write, backfill, run shadow reads until they match, flip reads to the target, then stop writing the source. With a version on every write and "apply only if newer" on both shards, both cases converge (S4 keeps v5 in the first, both end at v6 in the second). The version must come from one authority per key (the source's log position or a per-key counter), never from router clocks, which disagree (page 03, Part 3). Discord's migration backfilled with the original write timestamps (case study, step 3.3), and Uber's Mezzanine backfill and mirrored writes converged because both wrote each trip version under the same ref key (step 2.3).
Moving key by key
Valkey and Redis Cluster move a slot without freezing it, one key at a time:
| Stage | Source node | Target node | A client asking for a key |
|---|---|---|---|
| Start | Slot marked MIGRATING | Slot marked IMPORTING | |
Keys move with MIGRATE | Serves keys it still has; for a key it no longer has, replies ASK | Serves a request only if it is preceded by ASKING; otherwise replies MOVED back to the source | ASK: send only the next request to the target, without updating the map |
| A multi-key request whose keys are split between the two, or missing | Replies TRYAGAIN | Retry after a short delay | |
| Ownership changes | Replies MOVED | Owns the slot | MOVED: update the map for good |
It avoids a slot-wide pause, but a large key moves in one piece (ElastiCache doesn't migrate keys holding a large item during online resharding, and won't delete a shard holding an item over 256 MB), and multi-key requests on a moving slot keep retrying. Our fence is the whole-slot version: one short pause, one switch.
| Move method | Pause | What correctness needs | Bandwidth | Rollback |
|---|---|---|---|---|
| Log tail + fence (this Part) | Fence to catch-up: 13 ms here | A source slot table that refuses for good; a caught-up copy source | One copy plus the tail | Drop the target before the flip; the reverse move after |
| Dual writes + backfill | None, if versions converge | A version from one authority on every write; "apply only if newer" | Every write twice for the whole move | Before the read flip: stop writing the target; after: flip reads back while both are still written |
| Key by key (Valkey) | Per key, while it moves | ASK/MOVED handling in every client | One copy | Migrate the keys back |
| Freeze the slot, copy, switch | The whole copy: 10 s for 200 MB here | Nothing else | One copy | Unfreeze on the old owner |
What the move costs everyone else
The copy reads S2's replica's disk and uses the links between AZs while S2 keeps serving, so it is throttled below the spare I/O (Cassandra's stream_throughput_outbound defaults to 24 MiB/s for the same reason). When the hot slot is big and the cold ones small, it can be cheaper to move the cold slots away instead: p1, p5 and p13 (slots 11,819, 11,951 and 10,396) are three small copies that give the same isolation (event 24 in Part 12).
At scale. A 10 TB cluster going from 4 to 5 shards, as in Part 6, moves a fifth of its data: 2 TB. At 100 MB/s aggregate, throttled below the sources' spare I/O: 2,000,000 MB ÷ 100 MB/s = 20,000 s, about 5.6 hours. Buckets move in parallel, each with its own milliseconds-long pause. The move's length is set by bandwidth, and each bucket's pause by catch-up.
Engines that do it for you
| Engine | How it moves data | What to know |
|---|---|---|
| Vitess | VReplication copies and tails; VDiff compares; traffic is switched for reads and writes, and a reverse stream allows ReverseTraffic | A rollback path exists only while reverse replication runs |
| Citus | The shard rebalancer moves shards with logical replication while reads and writes continue, with a brief lock at the end; isolate_tenant_to_new_shard gives one tenant its own shard | The drill's big tenant, as a function call |
| Aurora PostgreSQL Limitless Database | Sharded tables are split by a hash of the shard key; routers hold the metadata. A shard split runs as an asynchronous job, and AWS states that finalizing the split "causes some downtime" | Plan splits outside peaks |
What to remember from Part 8
- Copy, tail, compare, fence, flip, and only then delete.
- The old owner refuses the moved slot for good, by checking its own slot table: some router always has an old map, and no timer can tell you how old.
- A move shares disks and links with live traffic: its length is set by bandwidth, each pause by catch-up.
Part 9. Across shards: queries and transactions
Everything so far touched one key and one shard. Some queries don't name a key, and some changes touch two. At 20:15 the app shows "Trending now": the top 10 posts by likes in the last hour, from every shard.
Trace: "trending now"
The router asks every shard (now five: S0 to S4) for its top 10 and merges them. Each reply is about 20 KB (labelled), and clients ask 1,000 times a second.
Synthesizing vector architecture diagram...
Snapshot S6. What to notice: the answer waits for the slowest of the five replies, and from AZ a four of the five replies cross an AZ boundary.
The slowest shard decides
A scatter-gather's latency is the slowest reply's, not the average. If each shard is slow (over 50 ms, labelled) 1% of the time, independently (a slow reply takes 50 to 200 ms, a normal one 3 to 7 ms, labelled), the chance that at least one of N is slow is:
| Shards N | 1 − 0.99^N | From the reference implementation |
|---|---|---|
| 1 | 1.00% | |
| 5 | 4.90% | Monte Carlo: 4.94% |
| 16 | 14.85% | |
| 100 | 63.4% | Monte Carlo: 63.36% |
At 5 shards in the reference implementation, a query's 99th-percentile latency was 168.5 ms against 51.0 ms for a single shard: the rare slow shard is now in every hundredth query's path. Two standard defences: a per-shard timeout that returns partial results (for a "trending" list, 4 of 5 shards is fine), and hedged requests (after a short delay, send a second request to a replica and take the first answer).
A popular query is a hot read
The same top 10, asked 1,000 times a second, is Part 4's hot read again. Compute it once a second in one job and cache it. The scatter-gather then runs once a second, not 1,000 times, and moves 5 × 20 KB = 100 KB a second.
The bytes matter too. EC2 bills data crossing AZs at 3,456 a month** at 4.15**. Salting multiplies the same bytes by S for every uncached read.
Local and global indexes
| # | Index | A query "trending in the last hour" | Cost |
|---|---|---|---|
| 25 | Local index (each shard indexes its own data) | Asks every shard | Scatter-gather: tail latency, N requests |
| 25a | Global index, partitioned by its own key, hour#<0-7> (salted, because one hour is a hot key) | 8 reads, not every shard | The index is updated asynchronously, so it lags the posts; a write updates two places |
Reports run on a copy
A platform-wide report, such as the drill's admin report over all tenants, shouldn't scatter across production shards at all. Stream every shard's changes to a warehouse or data lake and run it there, seconds to minutes behind: the wallet loop does it with S3 and Athena (step 2.1), and the hotel and payment loops do the same.
A platform-wide report needs every tenant's data, and the store is sharded by tenant. Where should it run?
Transactions that span shards
| # | Time | Event |
|---|---|---|
| 26 | 20:16:00 | nova follows alice. Two rows must change: following:{nova} (slot 4,778, on S1) and followers:{alice} (slot 749, on S0) |
| Option | How | Cost |
|---|---|---|
| Co-locate | Choose keys so both rows share a tag | Impossible here: an edge belongs to both users, and each user's rows must stay with that user |
| Two-phase commit | A coordinator prepares both shards, then commits | Locks held across round trips; a coordinator crash leaves both rows locked until it recovers. Keep it off the hot path |
| Saga with an in-flight state | Write one side, then the other, with a compensating step if the second fails | Readers can see the half-done state; the digital wallet loop uses an in_transit account for exactly this (step 2.2) |
| One write plus a derived index | Write following:{nova} only; a consumer of the change stream builds followers:{alice} | The reverse side lags by the stream's delay; the Change Streams & the Transactional Outbox loop primitive makes it reliable |
For a follow, the derived index is the natural fit: nobody needs alice's follower list to change in the same millisecond. "Who liked p9?" is the same shape: the liked table is keyed by user, so the per-post list is a derived index built from like-events, or a scatter.
What to remember from Part 9
- A query without the key asks every shard and waits for the slowest; if it's popular, cache its answer.
- Keep transactions inside one shard by choosing the key; otherwise a saga or a derived index, not 2PC on the hot path.
- Reports across shards run on a copy.
Part 10. End to end through the layers
Where is each decision actually taken? Follow one like on p9 from a phone whose request lands on a router in AZ b at 20:20, after the move, and one read.
Synthesizing vector architecture diagram...
What to notice: the like is acknowledged after step 1; step 2 is the router's periodic flush, which the shard checks against its slot table and epoch before it reaches storage.
| Hop | What it knows | What it absorbs | How it fails |
|---|---|---|---|
| Client | Nothing about shards | Nothing | Retries on timeout: the like needs an idempotency key (page 04) |
| Edge | The URL | Identical reads, collapsed into one fetch | Serves stale objects if configured to |
| Router | The tag and slot, its cached directory version, per-key limits | Hot reads (micro-cache, coalescing), hot writes (the 500 ms buffer), excess load (429) | A crash loses its unflushed buffer, rebuilt from durable events; an old map is bounced by MOVED |
liked and like-events | The user's key; the stream's partition key | Nothing: each like is one record | Unavailable: the like is refused. One write, so never half-done; the event follows from the change stream |
| Shard leader (S4) | Its slot table and epochs | Nothing | Refuses slots it doesn't own (MOVED, "retry"); throttles above its ceiling |
| Storage | The log and its positions | Nothing | Durability and replication: the Write-Ahead Log, fsync & Group Commit and Replication loop primitives |
A read of p9 stops at the router's micro-cache almost every time; about once a second per router it reaches S4 as one multi-key request (the post plus six sub-keys), and on the way back it refills the cache. Hot keys are cheapest to absorb near the client.
Sizing per AZ
Routers run in every AZ, so we size them to survive losing one. Peak traffic: 3,150 likes (15 × 50 + 2,400) plus 10,500 reads (15 × 100 + 9,000) = 13,650 requests a second. Each router takes 2,500 (labelled) and should run at 80% or less, so 2,000.
| Plan | Arithmetic, from the reference implementation | Normal load | After losing an AZ |
|---|---|---|---|
| Our story's 6 routers (2 per AZ) | 13,650 ÷ 6 = 2,275 each | 91% | 4 left: 3,412.5 each, 136.5%: overloaded |
| 9 routers (3 per AZ) | Enough to survive an AZ at 100%: 6 left × 2,500 = 15,000 ≥ 13,650 | 61% | 6 left: 2,275 each, 91% |
| 12 routers (4 per AZ), sized at 80% | 13,650 ÷ 2,000 = 6.8, so 7 must survive in two AZs: 4 per AZ, 3 AZs | 1,137.5 each, 45.5% | 8 left: 1,706 each, 68.3% |
Fleets are placed per AZ, so round up per AZ, and size for the loss of one.
What to remember from Part 10
- Every layer knows a little: the router the map, the shard its slots and epochs, storage its log.
- Hot keys are cheapest to absorb near the client.
- Size fleets per AZ, for the loss of one.
Part 11. On AWS
Every AWS data service has a unit, a way keys map to it, and a ceiling per unit. Knowing those three things tells you where your hot key will hurt. Stated only from AWS's public documentation and papers.
Managed services that use it
| Service | Its unit, and how keys map to it | Ceiling per unit | How it splits or grows | What AWS states |
|---|---|---|---|---|
| Amazon DynamoDB | Partition, chosen by "an internal hash function" of the partition key; one partition key's items are stored together, sorted by sort key | 3,000 read and 1,000 write units a second per partition, and per isolated item | Adds partitions transparently for throughput and storage; adaptive capacity raises a hot partition's share "automatically and instantly" and can isolate one hot item on its own partition; item collections split by sort key unless the sort key is monotonic or the table has an LSI; burst holds up to 300 s of unused capacity | Write sharding with random and calculated suffixes; GSI write sharding; on-demand handles double the previous peak and can throttle if that is exceeded within 30 minutes; Contributor Insights in two modes, with GSI write throttling missing from the throttled-items graph. The DynamoDB paper: a split for consumption takes minutes and doesn't help one hot item |
| Amazon Kinesis Data Streams | Shard; MD5 of the partition key → a 128-bit hash key → the shard owning that range | Writes: 1 MB/s or 1,000 records/s; reads: 2 MB/s and 5 GetRecords calls a second | SplitShard, MergeShards; UpdateShardCount at most double or half per call and 10 times per rolling 24 hours; parents go CLOSED, then EXPIRED, and are read before their children; on-demand mode scales shards for you | A key's records stay in order only within its lineage |
| Amazon MSK | Kafka partition; toPositive(murmur2(key)) % numPartitions | Broker-dependent | Add partitions; Kafka "does not currently support reducing the number of partitions", and it "will not attempt to automatically redistribute existing data" | Adding partitions changes which partition a key's new records go to; the order consequence is page 05's (Part 4) |
| Amazon ElastiCache (Valkey, Redis OSS, cluster mode) | 16,384 hash slots, CRC16(key) mod 16,384, hash tags | Node-dependent | Online resharding and rebalancing while the cluster serves, with some performance degradation; keys holding a large item aren't migrated, and a shard holding an item over 256 MB can't be removed | MOVED, ASK and TRYAGAIN replies, as in Part 8 |
| Amazon Aurora (Aurora PostgreSQL Limitless Database) | Sharded tables by a shard key, hash-based; collocated and reference tables; routers hold the metadata and manage distributed transactions | A shard can grow to 128 TiB | Shard splits run as an asynchronous job; finalizing a split "causes some downtime"; add routers for throughput | A managed version of the wallet loop's application-level sharding |
| Amazon OpenSearch Service | Primary shards: shard_num = hash(_routing) % num_primary_shards, where _routing defaults to the document _id and can be set per document | Node-dependent | The split index API: the target's primary shard count must be a multiple of the source's, bounded by number_of_routing_shards; the index must be made read-only first | Custom routing keeps a tenant's documents on one shard, and makes that tenant one hot shard |
| Amazon Keyspaces | Cassandra partition key | 3,000 read and 1,000 write units a second per partition | Splits partitions automatically | One hot row still meets the per-partition ceiling |
| Amazon S3 | Key prefixes | At least 3,500 PUT/COPY/POST/DELETE and 5,500 GET/HEAD requests a second per partitioned prefix | Scales gradually as load grows, returning 503 Slow Down meanwhile | A store that splits by load, like Part 6's range splits |
CloudFront request collapsing and Origin Shield do the hot-read half of this page at the edge: identical requests for one object become one fetch to the origin (Parts 4 and 10). They place no data, so they aren't a sharding service.
Running it yourself
| Option | What it is | Facts and sizing |
|---|---|---|
| Amazon EC2 with Amazon EBS (or instance-store NVMe) running Vitess, Citus, Cassandra or ScyllaDB, Kafka, or Valkey | You own the key, the directory, the slot tables, the moves, the throttling and the cutover | Cassandra: num_tokens 16 with allocate_tokens_for_local_replication_factor 3, and stream_throughput_outbound 24 MiB/s by default. Citus: logical-replication rebalancing with reads and writes continuing, isolate_tenant_to_new_shard, alter_distributed_table. Vitess: VReplication, VDiff, switching reads and then writes, and reverse replication for rollback |
| Sizing in words | Keep each shard at or under 70 to 80% of its measured ceiling at peak, after its share of salted keys; size routers and coalescers per AZ, rounded up, to survive an AZ loss (Part 10: 12 routers at 80%); throttle moves below spare I/O (2 TB at 100 MB/s ≈ 5.6 h) | Cross-AZ traffic on EC2 costs 3,456 a month uncached against $4.15 cached, for one query) |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like sharding | Why it isn't |
|---|---|---|
| Elastic Load Balancing (ALB, NLB) | "Spreads the load" | Spreads requests over stateless targets; it places no data |
| DynamoDB on-demand capacity | "No capacity planning, so no hot partitions" | Per-partition limits still apply to every key |
| RDS and Aurora read replicas | "More copies, more capacity" | Scale reads of the whole data set; every copy takes every write (page 06) |
| A bigger instance, or Aurora Serverless v2 scaling | "More capacity" | Vertical: one writer, and one hot row is still one lock |
| DAX | "A cache in front of a hot key" | Its query cache isn't invalidated by writes to the table: staleness you didn't choose |
| S3 Replication, DynamoDB global tables | "Copies spread over Regions" | Replication, not partitioning (page 06, and Multi-Region Failover, coming) |
What to remember from Part 11
- Know each service's unit and its ceiling: a DynamoDB partition, a Kinesis shard, a Kafka partition, a cluster slot.
- Managed splitting never goes below one key.
- You still choose the key, handle hot keys, and plan partition-count changes.
Part 12. What you've learned
Back to the celebrity post
Our table was 20% busy while one post took 33 times a partition's read ceiling and 10 times its write ceiling. In the tiny store, p9 put S2 at 255% of its write ceiling and 310% of its read ceiling. The pieces:
- One key lives in one piece (Part 1), and hashing spreads keys, not load (Part 2). Adding shards moved other keys, never
p9(Part 3). - Hot reads stopped at the routers: a 1-second micro-cache with single-flight cut 9,000 reads a second to 6, and coalescing routed by key cut uncacheable reads to 196 (Part 4).
- Hot writes were summed at the routers: 2,400 likes a second became 12 absolute-total flushes, safe to retry and rebuildable from durable events, because nothing else wrote
p9's slot per like (Part 5). Salting was the alternative, at 8 reads per count. - Growth moved exactly the slots we chose (Part 6); DynamoDB and Kinesis split for us, but never below one key (Part 7).
p9got its own shard while it was being liked: copy, tail, compare, fence, flip, with a 13 ms pause and a stale router bounced by the old owner's slot table (Part 8).- Queries across shards paid the slowest shard's latency, so the popular one was computed once a second (Part 9), and the routers were sized per AZ (Part 10).
- The costs: a second of staleness, 500 ms of count delay, a rebuild job, a directory, and slot tables that remember every move forever.
The whole story, event by event
From the reference implementation:
| # | Time | Event | Result |
|---|---|---|---|
| K1 | side | Three layouts, baseline | Range by minute: 800 writes/s on one shard (80%); hash by post: 200 each; directory by creator: 200 each, with nova's four posts together |
| K2 | side | Queries per layout | Counter: hash, directory; last minute: range; nova's posts: directory; anything else scatters |
| K3 | side | Key balance | Hash by post: exactly 4 posts per shard |
| K3a | side | Shard count (wallet) | 38,750 ÷ 2,800 = 13.8 → 14 → 16, for 64 buckets each |
| 1 | 20:00:00 | Baseline | 200 writes/s (20%), 400 reads/s (13%) per shard |
| 2 | 20:01:00 | nova promotes p9 | 2,400 likes/s, 9,000 reads/s by 20:02 |
| 3 | 20:02:00 | S2 melts | 2,550 w/s (255%), 9,300 r/s (310%); cluster 78.75% / 87.5%; p1, p5, p13 throttled |
| 4 | 20:02:05 | Detection | p9 = 94.1% of S2's writes, 96.8% of its reads |
| 5 | side | 4 → 8 shards | Shard 5 gets p9, p1, p5, p13: 2,550 / 9,300, unchanged |
| 6 | 20:02:10 | Micro-cache | 6 requests/s for p9 |
| 6a | side | Coalescing per router | 176.5 × 6 = 1,059 reads/s |
| 6b | side | Coalescing routed by key | 196 reads/s; 562.5 with one coalescer per AZ |
| 7 | 20:02:15.000 | Caption edit | Routers show it 0.187 to 0.779 s later; bound 1 s plus a query (3 s if refreshes fail) |
| 7a | side | Replica reads | Ceiling 3 × 3,000 = 9,000 < 9,300: 103% |
| 7b | side | 8 read copies | 1,125 reads/s each; every edit writes 8 |
| 8 | 20:02:20 | Salting, S = 8 | Two sub-keys per shard: S2 750 (75%), others 800 (80%) |
| 8a | side | Random salts; salt outside the tag | 96.15% chance of ≥ 3 on one shard (1,100 = 110%); likes:{p9}#3 all on slot 12,067 |
| 8b | side | Count reads | 72,000 sub-key reads/s uncached (600% per shard); 48 with the cache |
| 9 | 20:03:00 | Pre-aggregation | 12 flushes/s on slot 12,067; S2 at 162 writes/s |
| — | 20:03:05.410 | R2's retried flush | Ignored: seq 11 not larger |
| 10 | 20:03:07.300 | R4 crashes | 211 unflushed likes; sub-key stays 2,772 (seq 14); restart at 20:03:09.000 continues; the rebuild from like-events restores the 211 |
| 11 | side | Dedicated shard before the fixes | 2,400 w/s (240%), 9,000 r/s (300%); neighbours recover |
| K7 | day 2, copy | Add S5 | mod N: 80% move; ring with 16 tokens: 17.2% (median 19.9%), unevenly; fixed slots: 3,277 slots, 19.8% of keys, p13, p14, p15, p16 chosen; S5 40%, others 30% |
| K7a | day 2, copy | Ring, 1 token per shard | S5 takes 34.5%, all from S3 |
| K7b | day 2, copy | mod N doubling | 4 → 8: 50%; 12 → 16: 75% |
| K7c | day 2, copy | Split S2 by load | Cut at 12,067: p13, p1, p5 on one side, slot 12,067 alone; no finer split exists |
| D1 | side | DynamoDB item | 1,000 WCU served, 1,400/s throttled; SK = liked_at can't be split around |
| D2 | side | Hot GSI key nova | 2,550 writes/s on one index key; nova#0..7: 318.75 each |
| K4 | side | Kinesis, key post_id | 3 / 4 / 3 / 6 posts per shard; p9's shard 2,650 records/s, 1,650 rejected; its split child 2,500; random key 787.5 per shard |
| K5 | side | Packing 9 per record | 350 records/s; key can't be post_id |
| K6 | side | UpdateShardCount 4 → 8 | p9 on shard 6; parents CLOSED; read parents first |
| F1 | side | Kafka 4 → 8 partitions | 50.4% of keys' new records move; p9 stays in partition 2; old records stay; no reduction |
| 14 | 20:10:00.000 | Move starts | Directory v7; copy from a replica at 48,200 |
| 15 | 20:10:10.000 | Copy done | 200 MB at 20 MB/s = 10 s |
| 16 | 20:10:10.421 | Tail caught up | Backlog 1,619 entries, 121 for the slot |
| 17 | to 20:11:10.400 | Shadow reads | 358 compares; 16 caught S4 a few ms behind; all matched by position |
| 17a | 20:10:05.000 (branch) | S4 dies mid-copy | v7 unchanged, S2 never fenced; target dropped; restart from a fresh snapshot |
| 18 | 20:11:10.400 | Fence | S2: migrating, epoch 8, after 59,647; reads and writes get "retry" |
| 19 | 20:11:10.412 | Flip | S4 through 59,648 at .403; v8 committed .412; slot tables .413; 13 ms without writes; a probe flush 14 ms late |
| 19a | 20:11:10.405 (branch) | S4 dies before v8 | v9: S2 owner at epoch 9 at .502; resumes at .503; 103 ms pause |
| 20 | 20:11:10.418 | Push | Five routers on v8 |
| 21 | 20:11:10.500 | R5 on v7 | Flush 196,298 refused MOVED 12067 S4 epoch 8; applied at S4 at .504 |
| 21x | branch | No slot-table check | 42 flushes to the old copy until R5's poll at 20:11:31.200; a comment lost for good |
| 22 | 20:12:10.400 | Cleanup | S2 deletes its copy; answers MOVED forever |
| 23 | branch | Dual writes | Without versions: S4 left at v4, and v5 vs v6 diverge; with versions: both converge |
| 24 | branch | Move the cold slots | p1, p5, p13 (11,819, 11,951, 10,396): three small copies, same isolation |
| 25 | 20:15:00 | Trending now | 1 − 0.99⁵ = 4.90% slow; 63.4% at 100 shards; computed once a second: 100 KB/s |
| 25a | side | Global index hour#<0-7> | 8 reads instead of every shard; asynchronous |
| 26 | 20:16:00 | nova follows alice | following:{nova} on S1, followers:{alice} on S0: a derived index |
| 27 | 20:20:00 | One like end to end | Durable in liked and like-events, then OK; flushed to S4 within 500 ms; routers sized 12 at 80% |
The cheat card
| Topic | Remember |
|---|---|
| Routing | Key → tag → slot (CRC16 mod 16,384) → directory → shard; the shard checks its own slot table |
| Hash tags | Only the text in the first {...} is hashed; the salt goes inside the tag |
| Choosing a key | The most frequent query and the transaction boundary; never lead with time |
| Skew | Hashing fixes key skew, not load skew |
| Shard count | max(peak writes ÷ (ceiling × 0.7), data ÷ per-shard move size), rounded up; buckets ≫ shards |
| Hot key | One key's demand above its shard's ceiling; more shards never help it |
| Hot reads | Micro-cache with single-flight (stated staleness); coalescing routed by key: λ ÷ (1 + λt) |
| Hot writes | Salting (size after other traffic, check placement, reads × S) or pre-aggregation (absolute totals with a sequence; nothing else writes the key per event) |
| Dedicated shard | Protects neighbours; never more than one shard's ceiling |
| Adding capacity | mod N moves most keys; rings with many tokens about 1/N; fixed slots exactly what you choose |
| DynamoDB | 3,000 reads, 1,000 writes per partition and per item; adaptive capacity instant; split for consumption takes minutes, never one item; hot GSIs throttle the table |
| Kinesis | 1 MB/s or 1,000 records/s per shard, shared; a split never splits one key; read parents first |
| Live move | Copy from a caught-up copy, tail, compare, fence, flip, delete; the source answers MOVED forever |
| Dual writes | Versions from one authority per key, "apply only if newer" |
| Scatter-gather | Waits for the slowest: 1 − (1 − p)^N; cache popular ones; reports on a copy |
| Cross-shard change | Co-locate, or a saga, or a derived index; 2PC off the hot path |
| Sizing | Per AZ, rounded up, surviving the loss of one |
Failure checklist
- Does the key follow the most frequent query and the transaction boundary, and does no key lead with time?
- Is every shard's utilization alarmed against its own ceiling, not the cluster average?
- Can you name the hottest keys per shard within seconds (top-K at the routers, or Contributor Insights including the GSIs' most-accessed graphs)?
- Is there a per-key and per-tenant limit that answers
429before a shard melts? - Does every cache in front of a hot key have a stated staleness bound and single-flight on misses?
- Are salts inside the hash tag, placed or checked so they spread, and is the merged read cached?
- Does anything still write a pre-aggregated hot key once per event?
- Are flushes absolute totals with a sequence, and are the events durable before the "OK"?
- Does every shard check every request against its own slot table, and does the old owner refuse a moved slot for good?
- Is every live move throttled below spare I/O, copied from a caught-up source, and able to abort or roll back?
- Do readers of Kinesis or DynamoDB Streams finish a parent shard before its children, and is any Kafka partition increase planned as an operation?
- Are routers, coalescers and shards sized per AZ for the loss of one?
Think-first drills
Drill 1. A counter gets 6,000 writes a second. Shards take 1,000 writes a second, already carry 400 a second of other traffic, and must stay at or under 80%. How many salts, if you can check where each lands, and what does a read of the count cost?
Drill 2. 9,000 identical reads a second, 30 ms per query under load, 6 routers. How many reads reach the shard with coalescing on each router, and with coalescing routed by key? What does this say about coalescing when the shard slows down?
Drill 3. Size a new store: peak 20,000 writes a second, a shard takes 1,000 and should run at 70%; 6 TB of data, and you'll move or restore at most 250 GB per shard. How many shards and buckets? Then, growing from 32 to 40 shards, how many buckets move, and how many keys would hash mod N have moved?
Interview questions
| Question | Model answer |
|---|---|
| How do you choose a partition key, and what goes wrong with a time-based one? | List the most frequent queries and the transaction boundary, and pick the key they carry (user, tenant, entity), hashed into many fixed buckets mapped to shards by a directory. Check key skew and load skew. A key that leads with time sends every new write to the newest range: one hot shard, moving at every boundary. Bucket growing keys by time inside an entity key, or shard time keys into sub-keys. |
| One key gets 100× the traffic of the rest. Walk me through the fix for reads and for writes. | Find it per key, not per cluster. Reads: a short micro-cache with single-flight, and coalescing routed by key for reads that can't be stale; state the staleness. Writes: pre-aggregate at the writers and write absolute totals with a sequence, with the events durable elsewhere and nothing else writing the key per event; or salt, sized after each shard's other traffic, with the salt inside the tag and the merged read cached. Add a per-key limit as the guard, and a dedicated shard to protect neighbours. More shards don't help. |
| Consistent hashing vs fixed slots with a directory: when would you pick each? | A ring with many virtual nodes needs no central table, and adding a node moves about 1/N of the keys from many nodes; good for peer-to-peer membership (Dynamo-style stores, Uber's Ringpop). Fixed slots with a directory move exactly the slots you choose, by load, at the cost of keeping a highly available directory; most databases that move data deliberately do this. hash mod N is neither: it moves most keys. |
| How do you move a shard's data to a new node without downtime or lost writes? | Copy a snapshot from a caught-up copy, tail the source's log from that position, compare by position, then fence the source (it answers "retry" for the slot), wait until the target has applied the fence, commit the directory with a higher epoch, and flip both slot tables. The source answers MOVED for the slot forever, so a router with an old map is bounced however old its map is. Abort before the flip by dropping the target; roll back after it with the reverse move. |
| Your Kinesis stream throttles one shard. Why won't a split fix it? | A split halves the shard's hash-key range, but one partition key hashes to one point, so a hot key lands in one child with all its traffic. The limit is 1,000 records or 1 MB a second per shard, shared by every key on it. Use a random partition key if per-key order isn't needed, pre-aggregating downstream, or pack records if the key needn't name an entity. |
| How do you run a query or a transaction that spans shards? | A query without the key scatters to every shard and waits for the slowest (1 − (1 − p)^N), so use per-shard timeouts with partial results, cache popular answers, keep a global index for frequent lookups, and run reports on a copy. For changes, choose keys that keep transactions inside one shard; otherwise a saga with an in-flight state, or write one side and derive the other from the change stream; keep 2PC off the hot path. |
Where to go next
- Replication, Quorums & Read-Your-Writes: the copies inside each shard, and which copy may be a copy source.
- Leases, Fencing Tokens & Distributed Locks: epochs and fencing tokens, the rule behind the slot table.
- Idempotency & Effectively-Once Processing: why flushes and retries after
MOVEDmust be absolute or deduplicated. - Change Streams & the Transactional Outbox: the log tail that moves a slot, derived indexes, Kafka partition increases (Part 4) and stream lineage (Part 9).
- Event Time, Watermarks & Checkpoints: pre-aggregation inside a stream job (Part 3), and Flink's key groups when a job rescales (Part 9).
- Write-Ahead Log, fsync & Group Commit: what a shard's log position is.
- The Multi-Region Failover loop primitive (coming): cells, shuffle sharding and home Regions.
- Background: Primitive #08: Database sharding and partition keys, Primitive #01: Consistent hashing, Primitive #04: Distributed cache patterns and Primitive #10: Two-phase commit and sagas.
- Drills: The Product Page That Melted Redis (answered in Part 4), One Customer, One Shard, One Outage (Parts 5 and 8 for the tenant that grows 10×, Part 9 for the cross-tenant report) and The Ring Rebalance That Crushed Node 07 (Part 6).
- Loops: the key-value store (steps 1.3, 2.7, R2.8), the message queue (steps 1.4, 2.1, 2.7), the rate limiter (steps 2.1, 2.3, R2.6), the URL shortener (steps 2.1, 2.2, 2.5), the news feed (steps 2.1, 2.2), chat (steps 2.1, 2.2), the gaming leaderboard (steps 2.6, 3.1), notifications (step 2.2), S3-like storage (steps 1.1, 2.6), search autocomplete (steps 2.2, 3.2), the web crawler (steps 2.3, 3.1), metrics (steps 2.2, 3.1), the job scheduler (steps 1.1, 2.1), payments (step 3.1), the digital wallet (steps 2.1, 2.2, 2.3, R2.11, R3.2), the matching engine (steps 2.7, 3.6), hotel reservations (steps 2.1, 2.4, R3.9), the transactional outbox ledger (R1.11), the proximity service (steps 2.1, 2.5), maps and routing (steps 2.1, 3.1), nearby friends (steps 1.1, 2.3), ride-sharing dispatch (steps 1.1, 1.4, 2.6), YouTube (step 2.5), Google Drive (step 2.2), ad-click aggregation (steps 2.5, R2.6, R2.8), email (steps 2.4, 2.6), the mobile news feed (R2.5), mobile stock trading (steps 2.1, 2.3), and the Discord (steps 1.3, 2.1, 3.2, 3.3) and Uber (steps 1.1, 2.1, 2.3, 2.7) case studies.