Distributed systems
Part 1 of 6 · Consistent hashingConsistent Hashing: Rings, Virtual Nodes & Replica Placement
Modulo remaps ~all keys on membership change; consistent hashing remaps ~K/N via a clockwise hash ring. Vnodes fix skew and fan out failures; RF walks collect distinct physical nodes (topology-aware).
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What breaks in hash(key) % N when you add a node?
Answer
N changes, so almost every remainder changes and almost all keys move.
L2
Where does a key live on a consistent-hash ring?
Answer
On the first node at or clockwise from hash(key).
L3
One million keys, 10 nodes, add an 11th. About how many keys move on a ring?
Answer
About K/N, roughly one million divided by 11, near 91 thousand, not the full million.
L4
Why add virtual nodes instead of one token per machine?
Answer
One token skews load and dumps a neighbor's keys on failure. Virtual nodes spread both.
L5
How do you place RF replicas so they are not the same host?
Answer
Walk clockwise and skip tokens that belong to a physical node you already picked. Prefer another rack or zone.
L6
When would you pick rendezvous hashing or jump hash instead of a ring?
Answer
Rendezvous for a tiny weighted set with nothing to store. Jump for dense buckets and almost no memory. Neither places zone-aware replicas.
L7
The ring is balanced and a deploy still melts origin. What did hashing miss?
Answer
A hot key. Consistent hashing spreads keys, not requests for one key.
Failure modes
Token skew
One range owns far more keys, so one node runs hot while others idle.
Replica collapse
The next replica tokens all sit on the same host or rack, so one failure takes every copy.
Misconceptions
Consistent hashing stops a hot key from melting a node.
It balances key ranges. One viral key still lands on one owner.
Adding a node still moves about half the keys.
A ring moves about K/N, not half the keys.
Interviewer traps
Quoting modulo movement numbers for a ring.
Say about K/N for the ring and almost all keys for modulo.
Treating virtual nodes as extra physical replicas.
Virtual nodes are extra tokens on the same machine. The replica factor is a separate clockwise walk.
Design scenario
Same prompt for every reader.
Requirements
Design a product-catalog cache for 20 nodes. Adding or losing a node should move about K/N keys, not almost every key.
Traffic / scale
About 50k reads/s and 500 writes/s across 20 cache nodes.
Latency
Cached reads stay under a few milliseconds at p99.
Consistency
A key may be stale for a short TTL. The owner of a key stays put between membership changes.
Availability
Lose one node or one zone and still serve the key from another replica.
Failure assumptions
- One node or one zone can disappear.
- A single hot key can still overload its owner.
Constraints
- Do not use hash(key) % N as the placement rule.
Prompt
Sketch the ring, the vnode map, and the replica walk for this cache.
API
What does a client call to read a key, and what does it learn about the owners?
Data
What do you store on the ring, and how do virtual nodes map to physical hosts?
Architecture
Where does the ring live, and how does the replica walk pick distinct physical nodes?
- 1
membership change → almost every remainder flips
hash(key) % N
1M keys, 10 nodes, add an 11th: modulo moves ~1M. Cache deploys stampede origin. Do not ship this for stateful shards.
- 2
Winner for stateful storage: ring + vnodes + topology-aware RF
Same numbers move ~91k keys (~K/N). Vnodes cut skew and fan out failures. RF walks skip the same physical host and prefer other racks/AZs.
- ?
HRW / Jump / bounded-load — when they beat a ring
HRW: tiny N, weights, no ring to maintain. Jump: dense 0..N-1 buckets, almost no memory. Neither places AZ-aware replicas; neither fixes a viral key. See the sibling pages.
Why modulo sharding fails
hash(key) % N is fine until membership changes. Adding or removing one of N nodes changes N, so almost every remainder changes.
| Event | Keys remapped (modulo) | Keys remapped (consistent hash) |
|---|---|---|
| Add 1 of N nodes | ~100% (almost every key % N changes) | ~1/N |
| Remove 1 of N nodes | ~100% | ~1/N |
For caches this is catastrophic: a deploy that adds one cache host can stampede the origin. For databases it forces massive data movement and hot neighbors.
The hash ring
Place nodes and keys in the same circular hash space (typical teaching range: 0 .. 2^64-1). A key belongs to the first node clockwise.
Sequence
- 1
SortedTokenRing
Step 1 - Place nodes at hash(id) on the circle
- 2
ConsistentHashRouter → SortedTokenRing
addNode(A), addNode(B), addNode(C)
- 3
Client
Step 2 - Hash key, then walk clockwise to next node
- 4
Client → ConsistentHashRouter
get user:42
- 5
ConsistentHashRouter → SortedTokenRing
bisect hash(user:42)
- 6
SortedTokenRing → ConsistentHashRouter
owner = Node B
- 7
ConsistentHashRouter → Node B
GET user:42
- 8
SortedTokenRing
Step 3 - Add Node D: only keys in D's new arcs move
- 9
ConsistentHashRouter → SortedTokenRing
addNode(D)
- 10
SortedTokenRing
Step 4 - Failure: dead node's keys go to clockwise survivors
Lesson map
Consistent Hashing: Rings, Virtual Nodes & Replica Placement
Modulo remaps ~all keys on membership change; consistent hashing remaps ~K/N via a clockwise hash ring. Vnodes fix skew and fan out failures; RF walks collect distinct physical nodes (topology-aware).
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 client["Client"] router["ConsistentHashRouter is one of the participants this lesson's sequence actually names."] ring["SortedTokenRing"] nodea["Node A"] router -->|addNode(A),| ring client -->|get user:42| router router -->|bisect| ring ring -->|owner = Node B| router router -->|addNode(D)| ring
Ownership rule
- Hash the key into the ring space (64-bit with a good non-cryptographic hash — Murmur3, xxHash, truncated MD5 for demos).
- Find the smallest node token ≥ key hash (binary search on a sorted token array).
- If none, wrap to the first token (the ring is circular).
Complexity: lookup O(log V) where V is the virtual-node count. Membership updates rebuild or splice the sorted token list.
Virtual nodes (vnodes)
One physical point per server leaves large arcs and hot spots. Place each physical node at V positions:
- Identifiers like
nodeA#0,nodeA#1, …nodeA#(V-1) - Typical V: 100–256 (Cassandra historically used high
num_tokens; modern guidance often favors fewer, carefully allocated tokens) - Benefits:
- Load variance shrinks roughly like 1/√V
- Heterogeneous capacity: give stronger boxes more vnodes
- On failure, load fans out to many successors instead of one neighbor
Architecture
Without vnodes
- 1
Keys
- nextNode1 large arc
- nextNode2 small arc
- 2
Node1 large arc
- 3
Node2 small arc
With vnodes
- 4
Keys
- nextA-0
- nextB-7
- nextA-41
- nextC-12
- 5
A-0
- nextPhysical A
- 6
B-7
- nextPhysical B
- 7
A-41
- nextPhysical A
- 8
C-12
- nextPhysical C
- 9
Physical A
- 10
Physical B
- 11
Physical C
Replication on the ring
Primary = first distinct physical node clockwise. For RF = R, keep walking until you have R distinct physical nodes (skip other vnodes of the same host).
Topology-aware placement: prefer replicas in different racks / AZs so one failure domain does not take all copies. Zone configs, NetworkTopologyStrategy, and why R+W>RF is not enough if all copies share an AZ: Topology-aware replica placement.
Quorum example (Dynamo-style): RF=3, write quorum W=2, read quorum R=2 satisfies R + W > RF for strong quorum on the latest version (with vector clocks / versioning).
Failure paths
Decisions
- 1
Request for key K
- nextHash K onto ring
- 2
Hash K onto ring
- nextPrimary = next clockwise physical
- 3
Primary = next clockwise physical
- nextPrimary alive?
- ?
Primary alive?
- YesServe / coordinate
- NoWalk to next replica
- 5
Serve / coordinate
- nextMembership change: stream only ~K/N keys
- 6
Walk to next replica
- nextHinted handoff enabled?
- ?
Hinted handoff enabled?
- YesWrite hint on neighbor
- NoFail request or degrade
- 8
Write hint on neighbor
- nextReplay when primary returns
- 9
Fail request or degrade
- 10
Replay when primary returns
- 11
Membership change: stream only ~K/N keys
Hinted handoff is a durability convenience, not a substitute for RF. If the coordinator cannot gather a quorum, fail or explicitly degrade.
Membership views, streaming throttle, and vnode token moves: Vnode rebalancing & membership. A viral key still pins one primary — Hot keys & bounded loads.
Worked numbers
- 1M keys, 10 equal nodes → add 11th ≈ move ~91k keys (~1/11), not ~1M.
- Remove 1 of 10 → move ~1/10 of keys onto remaining owners (the dead node’s arcs).
- Without vnodes, one unlucky node can own 2–3× fair share; with V ≈ 150–200, imbalance typically drops under ~10–15% relative.
- RF=3 on N=10 physical nodes: each key has 3 owners; a single node death still leaves 2 copies if replicas were placed on distinct hosts.
Draw N=4 nodes, V=3 vnodes each, RF=3. Hash a key onto an arc owned by B#1. List the replica set (distinct physical). Then erase node B and show which keys move — and to how many successors.
Python ring (run it)
Self-contained router: add/remove nodes, lookup, and RF replica list. Run the demo: 3 → 4 nodes should move ~250 / 1000 keys, not ~1000.
Press Run. Snippets must be self-contained — no network, files, or native modules.
TypeScript ring (browser-safe hash)
Same contract. The playground has no Node crypto, so this uses a djb2 + avalanche mix as a stand-in for a 64-bit production hash (Murmur3 / xxHash).
Press Run. Snippets must be self-contained — no network, files, or native modules.
Interview Q&A
Why not hash(key) % N?
Answer
Almost every key remaps when N changes, because the divisor changed. Consistent hashing remaps ~K/N (more precisely ~K/(N+1) when adding one equal-weight node). For caches that is the difference between an origin stampede and a small rebalance.
What problem do virtual nodes solve?
Answer
Skew from random placement on a sparse ring (one unlucky node owns a huge arc). They also give weighted capacity (more vnodes → more share) and fan-out of failed-node load across many successors instead of one neighbor.
How do you place RF replicas?
Answer
Walk clockwise, collecting distinct physical nodes, until you have RF. Optionally skip the same rack/AZ for topology awareness. Never treat two vnodes of cache-a as two replicas.
What moves when you add a node? 1M keys, N=10, add an 11th.
Answer
Only keys whose clockwise owner becomes the new node — the arcs the new vnodes carve from former owners. Expect ~1,000,000 / 11 ≈ 91k keys for equal weights, not 1M. Quote ~K/N in the loop; compute K/(N+1) if they want the exact add case.
Do hot keys get fixed by consistent hashing?
Answer
No. The ring balances key cardinality, not request volume. One viral key still pins one primary (and its replicas). Mitigate with key salting / sharding the hot partition, local caching, or request coalescing. Bounded-load variants (steal to the next successor) help skewed assignment, not a single key’s QPS.
Alternatives to a vnode ring?
Answer
Rendezvous / HRW (highest random weight): score (hash(node, key)), pick max — simple, good rebalance, O(N) lookup unless approximated. Jump consistent hash (Google): O(1) memory, no vnodes, needs a compact 0..N-1 bucket set; awkward arbitrary weights and topology. Maglev: lookup tables for L4 LBs with bounded load.
Where is this used?
Answer
Cassandra / Dynamo-style stores, many Memcached clients, CDN / edge routing, sharded caches and partitioners. In interviews, name it whenever N stateful nodes change over time and reshuffling all keys would be fatal.
RF=3, W=2, R=2. Why is this a valid quorum — and what still fails?
Answer
R + W = 4 > RF=3, so a read quorum and a write quorum intersect on at least one replica (latest version, with versioning). It does not save you if all three replicas landed in one AZ, or if you counted vnodes of the same host as distinct, or if a hot key saturates the primary’s CPU.
What happens when a node fails — where do its keys go?
Answer
Lookups walk clockwise to the next live physical owner (and its replicas). With enough vnodes, that load fans out across many successors instead of dumping onto one neighbor. Hinted handoff can take writes for the down node and replay them on return. Membership changes should stream only the keys that moved — not a full reshape.
Why is vnode count a capacity knob?
Answer
More vnodes → smoother key share and finer weights (give a fat box more tokens). Too few: unlucky arcs and one-neighbor failure dump. Too many: bigger ring memory, slower updates, noisier membership. Pick enough that each physical node owns hundreds of tokens, then scale tokens with hardware, not with “more servers someday.”
Pitfalls
- Shipping without vnodes — small N almost always has severe imbalance.
- Counting vnodes as replicas — RF walk must skip the same physical node.
- Ignoring topology — all replicas in one AZ → correlated failure.
- Assuming perfect QPS balance — one viral key still melts one shard.
- Unbounded rebalance storms — throttle streaming when many nodes join/leave.
- Weak hash / few bits — collisions and clumping; use a solid 64-bit+ hash space.
- Client rings out of sync — gossip/config version skew sends traffic to the wrong owners; version the membership view.
Deep dive · Rendezvous, jump hashing, and Maglev — sibling lessons
Rendezvous (HRW) and Jump consistent hash are full studies in this cluster, not footnotes. HRW scores hash(key, node) and picks the max — O(N) lookup, ~1/N remap, easy weights, no ring. Jump maps onto 0..N-1 with almost no memory — you cannot hang arbitrary node ids or AZ constraints on it. Maglev is a huge lookup table that caps per-backend slots; prefer it for L4 / network load balancing (Maglev table), not for placing database replicas.
Default interview choice for stateful storage stays this page: ring + vnodes + topology-aware RF walk. Then say: “I would switch to HRW if N is dozens of caches; Jump if buckets are a dense integer range; bounded-load / PoTC if a viral key is the actual problem.” Read Rendezvous / HRW, Jump hash, and hot keys.