Distributed systems
Part 6 of 6 · Consistent hashingHot Keys & Bounded Loads: When Consistent Hashing Is Not Enough
Consistent hashing balances key cardinality, not QPS. Salt hot partitions, coalesce, or use bounded-load / power-of-two-choices so one viral key does not melt a shard.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
Every earlier page in this cluster makes key cardinality even. Rings, HRW, and jump hash all send user:42 to one home (plus replicas). If user:42 is a celebrity, that home eats 30% of QPS while its sibling sits idle. Vnodes, weights, and AZ walks do not help. This lesson is the overlay: detect heavy keys, split them, coalesce identical in-flight work, or steal assignments under a load cap.
Why bounded-load / two-choice wins for skew
| Tool | Balances | Cache / shard affinity | Movement on join | Handles one viral key? |
|---|---|---|---|---|
hash % N | Poor on churn | Sticky until N changes | ~all keys | No |
| Ring / HRW / jump | Key count ~1/N | Sticky home | ~K/N | No |
Salt key#i | QPS of that key across S shards | Split on purpose | Only the salted key | Yes, app-aware |
| Coalesce / singleflight | Waiters, not placement | Same home | None | Cuts duplicate work, not CPU of the one fetch |
| Power-of-two-choices | Request load ~log log N | Two homes | None until load changes | Spreads successive requests of the same key |
| Bounded load (1+ε) | Hard ceiling | Mostly sticky, spill when hot | Spill only over cap | Yes, with ε as the knob |
| Maglev table | Slot count per backend | Sticky until rebuild | Table rebuild | Caps connection share, not a single partition’s disk |
- 1
viral key → one primary
Ring / HRW / jump only
Perfect ~1/N key share. The celebrity key is still one vnode, one HRW winner, one jump bucket. Replicas help reads; the write/coordinator CPU is still concentrated.
- 2
application → S shards or 1 in-flight
Salt or coalesce
Salting is the honest database move (hot partition split). Coalescing is the honest cache move (one origin fetch). Neither needs a new hash family.
- ?
Winner for request routing: bounded-load two-choice
Two hashes, pick the lighter node unless both exceed (1+ε)·avg, then spill. Maglev wins when you can afford a permutation table and want O(1) LB lookup. Do not use Maglev to place RF=3 across AZs.
What “hot” means
| Symptom | Cause | If you only add nodes |
|---|---|---|
| One shard’s p99 explodes | One key (or tiny key set) dominates QPS | That key moves with its home; the new node does not siphon QPS unless the key remaps |
| CPU / lock / WAL pinned | Hot row inside an otherwise quiet partition | Extra replicas may help reads; writes still serialize |
| Cache hit ratio looks fine, origin melts | Many waiters on the same miss | Need collapsing, not more cache nodes |
| After a join, one node still hotter | Load metric ≠ key count | Classic bounded-load use case |
Consistent hashing’s ~K/N movement can make a hot key jump during a scale event and melt a different box. Rebalancing should throttle, but it will not split the key.
Power-of-two-choices (PoTC)
- Hash the key with two independent seeds (or two hash families) onto the node list / ring.
- Read a load metric on both candidates (in-flight requests, CPU, bytes, connections).
- Send the request to the lower load.
Why it works: throwing balls into bins uniformly gives max load ~ log N / log log N. Two choices, pick the lighter, drops the max to ~log log N with high probability. Independence of the two hashes matters; h and h+1 modulo N are not two families.
Affinity cost: the same key may alternate between two nodes as their loads seesaw. Caches then miss more. Databases that must own a key (single-primary partitions) cannot PoTC per request — they have to split the key instead.
Bounded-load consistent hashing
Mirrokni, Thorup, Zadimoghaddam: start from a consistent-hash home, but a node may hold at most
(1 + ε) · (total_load / N)
If the home is over cap, walk to the next candidate (the second choice, then clockwise, then global min). ε is the knob:
| ε | Ceiling | Affinity |
|---|---|---|
| 0 | Perfect average (often impossible without constant stealing) | Poor |
| 0.1–0.2 | Typical teaching range | Most keys stay home |
| Large | Almost never spills | Hot key stays stuck |
On node failure, redistribute that node’s load through the same rule so the cap still holds. Stale instantaneous counters lie under bursts — use a short EWMA.
Decisions
- 1
Incoming request
- nextHash key with two seeds
- 2
Hash key with two seeds
- nextCandidate node1 load L1
- nextCandidate node2 load L2
- 3
Candidate node1 load L1
- nextEither load at most (1+eps) times avg?
- 4
Candidate node2 load L2
- nextEither load at most (1+eps) times avg?
- ?
Either load at most (1+eps) times avg?
- YesAssign to lighter admissible node
- NoSpill to next least-loaded node
- 6
Assign to lighter admissible node
- nextServe and increment load
- 7
Spill to next least-loaded node
- nextServe and increment load
- 8
Serve and increment load
Lesson map
Hot Keys & Bounded Loads: When Consistent Hashing Is Not Enough
Consistent hashing balances key cardinality, not QPS. Salt hot partitions, coalesce, or use bounded-load / power-of-two-choices so one viral key does not melt a shard.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["Incoming request"] b["Hash key with two seeds"] c["Candidate node1 load L1"] d["Candidate node2 load L2"] a -->|Incoming request to Hash key with two seeds| b b -->|Hash key with two seeds| c b -->|Hash key with two seeds| d
Step labels: (1) two hashes, (2) read loads, (3) apply the (1+ε) cap, (4) pick or spill, (5) serve. Sticky consistent hashing stops at step 1 with a single home and no cap.
If both candidates are the same node (hash collision of two seeds), bump the second seed — otherwise you silently degraded to one-choice.
Maglev, briefly
Google Maglev fills a large table (often 65537 slots) with a permutation per backend. Lookup is table[hash % M]. Each backend gets a bounded number of slots, so a small backend cannot be assigned a huge fraction of the table. Membership change rebuilds the table (cost O(N · M) naively). Prefer Maglev for software L4/L7. Prefer a vnode ring plus this page’s salting for stateful shards.
Count-Min Sketch in one whiteboard
Count-Min is a small grid: d hash rows, w columns. Each event key increments table[i][ hash_i(key) % w ] for i in 1..d. The query is the min of those d cells — an overestimate, never an underestimate (in the basic version).
Pick w and d so that if a key is above threshold T (say 1% of QPS), you see it with high probability and the noise floor is below T. You do not store the key itself until it crosses T (keep a sidecar heap of current heavy hitters). Decay: multiply the table by α < 1 every window, or use sliding CMS, so last week’s Super Bowl ad does not stay salted.
Once user:hot is in the heap:
- Cache: coalesce + a process-local LRU; maybe replicate the object to all edges.
- KV / DB: emit a split plan — create
user:hot#0..#S-1, dual-write for one generation, then read scatter-gather. That is a targeted membership move of one key’s bytes, not ~K/N of the cluster. - Stateless router: flip that key onto bounded-load two-choice until the sketch cools.
Maglev table (why it is not a ring)
Maglev builds an array M[0..65536] (prime length). Each backend has a permutation of slots. The fill loop walks permutations and claims empty slots until every backend has about M / N entries, capped so a slow backend cannot own half the table. Lookup: M[hash(packet) % 65537].
| Property | Maglev | Bounded-load two-choice | Vnode ring |
|---|---|---|---|
| Lookup | O(1) array | Two hashes + load RPC/gauge | O(log V) |
| Bound | Slots per backend | (1+ε)·avg requests | None on QPS |
| State | ~256 KB+ table | Node list + loads | Token list |
| Rebuild | On membership | None | Token splice |
| AZ RF | No | No | Yes, with a walk |
Use Maglev in front of stateless or connection-oriented backends. Use salting + a ring for partitioned state.
Detecting hot keys
You will not store a counter per user id. The Count-Min section above is the data structure; this is the control loop. Optionally decay counts so yesterday’s viral clip does not stay salted forever. Detection feeds salting: once user:42 exceeds Q, write user:42#0 .. #S-1 and have the app pick a salt (random for writes that commute, or a sub-key). Alerting on “shard CPU” without a sketch leaves you guessing which key.
Salting and coalescing
- Salt / split:
order:9999#0,order:9999#1, … map to different homes. Reads may scatter-gather. This is how you fix a hot partition in a store that cannot change primary per request. - Local cache in front of the shard: the celebrity key is still one origin, but most QPS never arrives.
- Coalesce: N waiters for the same cache key → one fetch (singleflight). Complementary to hashing; it does not move the key.
Topology still applies: salted pieces should not all land in one AZ if each piece is RF=3.
Playground — naive vs two-choice vs bounded load
In-memory router. 30% of traffic is one key. Logs per-node load for sticky hash, PoTC, and (1+ε) spill.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Five cache nodes, 10k QPS, 3k QPS from user:hot. Draw sticky hash: one node at ~3.4k, others ~1.4k. Draw two-choice: the hot key’s two candidates share the 3k. Draw salt user:hot#0..#3. Then lower ε until the bounded-load router must spill to a third node — that is the affinity cost.
Interview Q&A
Does consistent hashing fix hot keys?
Answer
No. It equalizes how many keys each node owns, not how hot those keys are. One viral key still pins one primary (and its replica set). Quote this even after a perfect vnode diagram.
What is power-of-two-choices, and why does it crush imbalance?
Answer
Hash each request to two independent candidates, send to the lighter. Max bin load falls from ~log N / log log N to ~log log N. The two hashes must not be a trivial shift of each other.
How does bounded-load hashing differ from a plain ring?
Answer
The ring still proposes a home, but a node may not exceed (1+ε)·average. If both (or the home) are over cap, spill. Locality mostly holds; the cap is the guarantee. Plain rings have no QPS ceiling.
When do you prefer Maglev over PoTC?
Answer
L4/L7 with millions of lookups/s, need O(1) table lookup, can rebuild a large permutation table on membership change, want bounded slots per backend. Not for vnode RF or AZ placement.
How do you detect hot keys in a stream without a counter per key?
Answer
Count-Min Sketch or heavy-hitters with a threshold and optional decay. Sample if even the sketch is too hot. Detection must drive a split (salt) or a routing change, not just a dashboard.
What does increasing ε do?
Answer
Larger ε: fewer spills, more affinity, weaker ceiling — the celebrity key stays home longer. Smaller ε: tighter balance, more stealing, more cache misses / primary ping-pong. Tune from measured skew, not from 0.1 folklore.
Node failure under a load cap. What do you do?
Answer
Remove the node from the candidate set, reassign its in-flight and sticky keys through the same two-choice / bounded rule so survivors stay ≤ (1+ε)·new average. A full rebalance still streams data if the node held state.
Salt vs two-choice vs coalesce — pick one sentence each.
Answer
Salt: split one logical key into S sticky shards (databases). Two-choice: send successive requests of the same key to the lighter of two homes (stateless / cache). Coalesce: many waiters, one in-flight fetch (thundering herd). You often want salt and coalesce.
30% of QPS is one key, N=5 sticky caches. Rough loads?
Answer
One node gets ~30% + 14% of the remaining 70% if that remainder is fair, i.e. about 44% of cluster QPS on one box (0.30 + 0.70/5), others ~14%. Two-choice on that key shares the 30% across two nodes. Salting with S=5 can spread the 30% across all five.
Why can PoTC hurt a write-heavy KV store?
Answer
If each request of user:hot may land on a different primary, you lose single-key serializability unless those nodes are merely caches. Stores with one primary per key must split the key (salting / new partition) rather than bounce the primary per request.
Pitfalls
- “We use consistent hashing, so load is even.” Cardinality ≠ QPS.
- Stale load metrics — bursty traffic fools instantaneous counters; EWMA or in-flight gauges.
- Correlated hashes — two-choice collapses to one-choice.
- ε too large — the hot key never spills.
- ε too small on a stateful store — constant primary migration, worse than the hot key.
- Maglev for database RF — wrong abstraction; no AZ walk.
- Salting without a scatter-gather read path — you split writes and broke reads.
- No detection — you find the hot key in an incident review, not in a sketch.
- Coalescing only at one tier — five edges still mean five origin fetches unless each tier singleflights.
Deep dive · ε, EWMA, and when not to steal
Bounded-load papers analyze assignment of balls to bins with a cap. In production the “ball” might be a connection, a cache key, or a stored shard. Stealing a stored shard is a membership move: you must stream bytes and bump a generation. Stealing a stateless request is free. Do not mix those sentences.
Load signals: in-flight RPCs beat CPU percent (which lags). Exponentially weighted moving averages stop a 2 ms spike from spilling half the keyspace. For caches, a small local LRU of the celebrity key in every app task can beat any cluster-level steal.
If the hot object is read-mostly, topology-aware extra read replicas in more AZs plus an app cache may be enough. If it is write-heavy, salt or a dedicated queue for that key is the honest design.