Concept lesson · Foundations
Consistent hashing and virtual nodes
Start here
Definition
Consistent hashing assigns keys to owners so that adding or removing an owner changes only a limited portion of existing assignments. In the ring form, both keys and owner positions are hashed into one circular space, and a key belongs to its first clockwise owner.
Why it matters: The rule hash(key) mod N remaps many keys when N changes. A ring limits movement during cache expansion or shard membership changes, reducing cold misses and migration work.
Keys go clockwise to the first node token. Adding D at 40 moves (20,40] from B to D; other ranges keep their owners.
Read the diagram step by step
- Tokens are positions on a hash space, not geographic servers.
- Initially A=20, B=50 and C=80. Key 35 belongs to B.
- Adding D=40 transfers only keys in (20,40] from B to D. Key 45 still belongs to B.
- Virtual nodes improve distribution. A hot key can remain hot, and replication is a separate policy.
Worked example
On a 0–99 ring, owners A20, B50, and C80 place hash 35 at B50. Add D40: hash 35 moves to D40, while hash 45 stays at B50. Only the interval (20,40] changes owner.
Key takeaways
- First clockwise owner determines placement; the rule wraps past the largest token.
- Virtual nodes improve placement balance, but do not split one hot key.
- Consistent hashing is about placement stability, not CAP consistency or automatic migration.
You will learn to
- Map keys to clockwise owners, including wraparound, using actual numbers.
- Calculate which keys move on addition/removal and qualify expected movement.
- Explain virtual nodes, skew, and a safe migration without confusing placement with replication.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Data partitioning and sharding · Caching: cache hits, misses, write policies and invalidation
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01What is consistent hashing, and why not modulo N?
The common ring version hashes both keys and machine positions into the same circular number space. A key belongs to the next machine position clockwise. A virtual node, or token, is an additional ring position assigned to a physical machine, not another server. We will compute the placement before discussing migration and balance.
A distributed hash table associates keys with values and uses a deterministic rule to locate their owners. In a metadata-cache example, keys P12 and P35 identify records; a hash function maps each key to a numeric placement position. Placement stability determines how much cached or durable data must move when membership changes.
One cache is easy to address but eventually runs out of memory or throughput. With three caches, the application must decide where P35 lives. A common first rule is hash(key) mod 3. Everyone can calculate the same owner without a lookup table for every object.
Here mod means the remainder after integer division. Number the three destinations 0, 1 and 2; dividing a key’s hash by 3 produces one of those remainders, which selects its destination. All clients using the same hash and machine numbering therefore agree where to send that key.
The difficulty appears when a fourth machine joins. Changing the rule to mod 4 moves many keys. For a numeric hash of 35, the remainder changes from 2 to 3; for 12, it stays 0. Some mappings remain, but widespread movement can create cache misses or durable-data migration. We want a placement rule that changes fewer existing assignments when capacity changes.
02The hash ring: clockwise ownership with five keys
For a small worked example, let hashes range from 0 through 99. Connect 99 back to 0 to form a circle. Put cache A at position 20, B at 50, and C at 80. Real systems use a much larger space; the tiny range lets us compute every step by hand.
Node labels include their hash positions. The leader line locates key 55; the colored clockwise arc ends at its successor. Compare before and after adding D60.
Remember: Only the new owner’s predecessor interval moves.
Read the diagram
- Before the change, A20, B50 and C80 own the ring. Key 55 reaches C80 clockwise.
- After D60 joins, key 55 reaches D60 first.
- Only keys in (50, 60] move from C to D; other ownership remains unchanged.
Try from memoryWould key 65 also move to D60?
No. Clockwise from 65, the next owner is still C80. D60 takes only (50,60].
To locate a key, hash it and move clockwise until reaching the first cache position, including an exact match. That cache owns the key under our convention. P35 hashes to 35, so it reaches B50. P90 reaches the end of the range, wraps through zero, and reaches A20.
| Key | Hash | First clockwise position | Initial owner |
|---|---|---|---|
| P12 | 12 | 20 | A |
| P35 | 35 | 50 | B |
| P45 | 45 | 50 | B |
| P65 | 65 | 80 | C |
| P90 | 90 | 20 after wraparound | A |
Each position owns the interval after its predecessor and through itself. B therefore owns (20,50]: positions greater than 20 and less than or equal to 50. The round bracket excludes 20; the square bracket includes 50.
Hash collisions are expected in a placement space: two different object keys may map to the same number and therefore the same machine. Store and compare their full keys so they remain different records. Consistent hashing chooses an owner; it does not make a key unique. Clients must also use the same hash function, key encoding, token order, and membership version to calculate the same owner.
03Adding and removing a node: which keys move?
Add D at position 40. It becomes the first clockwise owner for hashes in (20,40]. B’s old interval splits: D takes (20,40], while B keeps (40,50]. Key P35 moves from B to D; P45 stays with B. Other intervals are unchanged.
Next remove B. Its remaining interval moves to the next position, C80. P45 now moves to C. This removal does not require moving P12, P35, P65, or P90.
| Key | Before addition | After adding D40 | After removing B50 |
|---|---|---|---|
| P12 | A | A | A |
| P35 | B | D | D |
| P45 | B | B | C |
| P65 | C | C | C |
| P90 | A | A | A |
Consistent hashing tries to preserve existing assignments where membership change does not require a new owner. The ring is a placement mechanism, not a promise that the new machine already contains the object. We still need to move or rebuild data and coordinate routing.
A20, D40, B50 and C80 sit at their numeric positions on the 0–99 ring. P35 lies between A20 and D40, so adding D moves P35 from B to D. P12, P45, P65 and P90 keep their owners, including P90 wrapping through zero to A20.
Read the key assignments
- P12 hashes to 12: A20 before adding D40; A20 afterward.
- P35 hashes to 35: B50 before adding D40; D40 afterward.
- P45 hashes to 45: B50 before adding D40; B50 afterward.
- P65 hashes to 65: C80 before adding D40; C80 afterward.
- P90 hashes to 90: A20 before adding D40; A20 afterward.
04Virtual nodes, balance, and physical failure domains
Our original intervals are unequal: A owns the wraparound interval (80,20], B owns (20,50], and C owns (50,80]. With uniformly distributed hashes, A owns about 40% of the space while B and C own about 30% each. Randomly choosing one position per machine can produce even larger imbalances.
Virtual nodes assign several positions to each physical machine. For example, A can own tokens A1 and A2 in separate parts of the ring. It then receives several smaller intervals rather than one possibly large interval. More well-distributed positions tend to smooth random imbalance and permit capacity-aware allocation.
To make virtual nodes concrete, use a separate six-token example with A at 10 and 60, B at 30 and 80, and C at 45 and 95. A owns (95,10] and (45,60]: 15 + 15 = 30 positions. B owns two 20-position intervals, totaling 40; C owns two 15-position intervals, totaling 30. Two tokens per host do not guarantee perfect balance. The benefit appears statistically or through deliberate token allocation across many smaller ranges.
A1 is server A’s token at hash position 10; A2 is its token at 60. Matching letters and colors group the tokens by physical server. The colored arcs show primary ownership. Two tokens per server still give unequal 30%, 40%, and 30% shares.
Remember: Several ring positions can point to one physical server.
Read the diagram
- This is the separate six-token example. Clockwise positions are A1 at 10, B1 at 30, C1 at 45, A2 at 60, B2 at 80, and C2 at 95.
- A1 and A2 belong to physical server A. Their primary ranges are (95,10] and (45,60], totaling 30 of the 100 hash positions. The first range wraps through zero.
- B1 and B2 belong to server B. Their ranges are (10,30] and (60,80], totaling 40 positions. C1 and C2 belong to server C and own (30,45] and (80,95], totaling 30 positions.
- A key hashing to 5 reaches token A1 at 10; a key hashing to 55 reaches token A2 at 60. Both keys are assigned to the same physical server A.
- Two tokens per server still produce unequal 30%, 40%, and 30% shares in this example. These percentages measure hash-space ownership, not necessarily bytes or request traffic.
- The diagram shows primary ownership. A1 and A2 share one physical failure domain; extra tokens do not create replicas. Replication must select other physical owners and appropriate failure domains.
Try from memoryIf physical server A fails, does its other token keep either key available?
No. A1 and A2 are positions assigned to the same server, so both lose that server together. Availability would require a usable replica on another physical server and a recovery protocol.
A production example is Cassandra's token-based placement: multiple tokens may belong to one node, while replica selection must skip duplicate physical owners. Increasing token count also adds placement metadata and more ranges to manage; choose it from operational needs rather than assuming the largest possible value is best.
A ring is not the only way to keep most assignments stable when membership changes. Another approach ranks the eligible machines separately for each key. A newly added machine takes that key only if it outranks the existing winner, avoiding the need for token positions.
Rendezvous hashing, also called highest-random-weight hashing, is another placement algorithm. Compute a deterministic score hash(key, nodeId) for each eligible node and choose the highest, using a stable tie-breaker. With illustrative scores A=.31, B=.86 and C=.54, the key belongs to B. Adding D with .70 leaves it on B; adding D with .93 moves it to D. Removing a node changes only keys that selected it. All routers need the same membership and scoring rules.
Unlike a token ring, the simple implementation evaluates all N nodes per lookup. It avoids virtual-node metadata but pays O(N) scoring cost, meaning the number of scores grows in proportion to the number of nodes; optimized variants and weighting require their own analysis. Neither placement method fixes a single hot key or performs safe data migration.
Interview check: Does adding a node move every key? No. A key moves only if the new node outranks its previous owner; moving durable bytes and changing write authority are separate steps.
05Estimate movement and recognize hot-key limits
Assume 1.2 million equal-sized metadata objects, equal-capacity machines, and balanced placement. Adding a fourth machine to three should move about one quarter of the keys, roughly 300,000, to the new machine on average. At an assumed 500 bytes per object, that is about 150 MB of payload before indexes, protocol overhead, or redundant copies.
More generally, adding one machine to N existing balanced owners moves an expected fraction near 1/(N+1); removing one of N owners moves near 1/N. These are distribution-based estimates. Our fixed D40 example takes a 20-position interval, not exactly 25% of the ring.
Fixed logical buckets offer another way to separate keys from physical machines. A bucket is a stable group of keys; a routing map records which machine currently owns each group. Changing that map can move selected groups without changing every key’s grouping rule.
| Placement choice | Useful when | Main resizing cost |
|---|---|---|
| Direct hash modulo machine count | Membership is fixed or remapping is cheap | Changing the divisor remaps many unrelated keys |
| Fixed logical buckets plus an owner map | Explicit migration batches and simple routing are useful | Maintain and distribute the bucket-to-machine map |
| Consistent-hash ring with tokens | Membership changes and limited reassignment matter | Maintain agreed membership, balance ranges, and migrate/refill them |
A logical bucket is a stable group of keys, such as bucket 17 of 1,024, that a routing map assigns to a physical machine. Moving bucket 17 changes its physical host without changing the hash modulus for every key. A ring is one good placement strategy, not a prerequisite for every sharded system.
06Data migration: copying, catch-up, and routing cutover
For an ordinary cache, D can start empty and fetch P35 from the authoritative database on a miss. But a sudden transfer of many hot keys can overwhelm that database. Warm selected keys, limit concurrent refills, and keep origin protection active during the change.
For durable storage, keep B’s copy until D is ready. Copy a consistent snapshot of the moving range, record and apply updates made during copying, then verify D’s data before making it the writer. Give the routing change a version so clients can detect old routes. During the move, keep one writer or use a protocol that explicitly coordinates the handover.
Suppose a write updates P35 while copying occurs. The destination must receive the newer version before it becomes authoritative, or the old owner must forward/reject according to the migration protocol. A stale client sending to B needs a safe redirect or forwarding path. Retain rollback information until verification completes. Ownership math says where P35 belongs; it does not implement this data-transfer protocol.
07Interview answer: draw the ring and explain the tradeoff
Interviewer: “What does consistent hashing solve when we add a cache?”
Candidate: “It reduces placement changes. On our 0–99 ring, P35 belongs to B50. Adding D40 transfers only the interval (20,40], so P35 moves to D while P45 stays with B. That can avoid the broad remapping from changing a modulo divisor.
“I would add virtual positions to improve placement balance, but I would still check object sizes and hot-key traffic. For a cache I need controlled refill; for durable storage I need snapshot transfer, concurrent-update catch-up, and safe routing cutover. The ring also does not provide read consistency or replication automatically.”
This explanation can be replayed on a whiteboard with five keys. It shows why the technique helps, when its balancing assumptions fail, and which essential migration decisions remain outside the hashing algorithm.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
What is consistent hashing? Draw a ring and explain why adding a node moves fewer keys than changing a modulo divisor.
Reveal a model answer
Consistent hashing is a placement scheme that limits remapping when owners join or leave. Draw a ring numbered 0–99 with A at 20, B at 50, and C at 80. Hash a key and choose the first clockwise owner, wrapping at 99. Hash 35 belongs to B50; hash 90 wraps to A20.
Add D40: it takes only (20,40] from B, so hash 35 moves to D while hash 45 stays at B. Changing hash(key) mod 3 to mod 4 would change many unrelated assignments. In balanced equal-capacity placement, adding one to N owners moves about 1/(N+1) of keys on average; this particular D40 interval covers 20% of our toy ring. Virtual nodes improve balance, but data still needs migration or cache refill, and one hot key remains a separate problem.
Interviewer follow-up
What happens for hash 90?
Reveal the follow-up answer
“It wraps through 99 and 0 to A20. That wraparound interval is part of A’s ownership.”
What the answer must demonstrate: Demonstrate the rule with actual positions.
On a 0–99 ring with A20, B50, C80 and keys at 12, 35, 45, 65, 90, which keys move when D40 joins?
Reveal a model answer
“Only P35 moves in our five-key sample. D takes (20,40] from B; P45 is outside that interval and stays with B. A and C keep their existing intervals. I would show the interval, not claim that every key moves to a new server.”
Interviewer follow-up
Does moving the sample key at 35 imply exactly one quarter of all keys moved?
Reveal the follow-up answer
“No. In this fixed ring D40 receives (20,40], which is 20 of 100 positions. An expected 25% movement requires four balanced owners and suitable hash-distribution assumptions; a five-key sample need not match either fraction.”
What the answer must demonstrate: Keep a concrete trace distinct from a statistical estimate.
On a clockwise ring with A20, D40, B50, C80, which owner receives B50’s interval when B is removed?
Reveal a model answer
“B’s remaining interval (40,50] passes to C80, the next clockwise owner. P45 moves to C. P35 stays with D. For durable data I must also ensure C obtains the required current state; the placement calculation does not transfer bytes.”
Interviewer follow-up
What if B fails before a copy is made?
Reveal the follow-up answer
“Recovery needs another durable replica or retained history. A ring alone cannot reconstruct missing data.”
What the answer must demonstrate: Placement and durability are separate responsibilities.
Why not just change hash(key) mod 3 to mod 4?
Reveal a model answer
“That changes many assignments at once, even though most existing machines are still healthy. Hash 35 changes remainder from 2 to 3, while 12 happens to stay at 0. Broad remapping can create expensive migration or cache misses; consistent hashing limits the affected ranges.”
Interviewer follow-up
Does modulo become impossible to use?
Reveal the follow-up answer
“No. It is simple for fixed membership or when managed logical buckets absorb physical changes. The issue is the resizing consequence.”
What the answer must demonstrate: Avoid claiming every modulo mapping necessarily changes.
What do virtual nodes improve?
Reveal a model answer
“They give one physical host several separated ring positions, so it owns multiple smaller intervals. With a suitable distribution, this reduces random placement imbalance and can represent differing capacities. It adds token metadata and migration units; it does not create more independent machines.”
Interviewer follow-up
How do you place three replicas when several consecutive virtual tokens belong to one physical host?
What the answer must demonstrate: Count physical failure domains for replication.
One key P35 receives half of all reads. Will more virtual nodes split that hot key?
Reveal a model answer
“No. The same key still maps to one primary owner under this rule. I would consider read replication, caching, or request coalescing, while defining update and freshness behavior. Virtual positions improve distribution across many keys rather than splitting one indivisible key’s traffic.”
Interviewer follow-up
What other imbalance should you measure?
Reveal the follow-up answer
“Bytes per object. Equal key counts can hide one owner holding much larger values and exhausting storage first.”
What the answer must demonstrate: Key count, bytes, and traffic are different load measures.
How many of 1.2 million keys move when three balanced owners become four?
Reveal a model answer
“The expected share for the new equal-capacity owner is about one quarter, or 300,000 keys. I would label the balance and distribution assumptions. At 500 bytes each that is about 150 MB of payload before overhead, which helps estimate a controlled transfer.”
Interviewer follow-up
Why can a particular node insertion move a different fraction than that expectation?
Reveal the follow-up answer
“The estimate assumes balanced placements and a suitable key distribution. On a 0–99 ring, adding D40 between A20 and B50 moves (20,40], only 20% of that fixed space.”
What the answer must demonstrate: Qualify both arithmetic and assumptions.
A write updates P35 while its ownership moves from B to D. What must the migration protocol guarantee?
Reveal a model answer
“D needs a snapshot and the updates committed while that snapshot is copied. I would catch up, verify, and atomically change the authoritative routing generation under the migration protocol. B must forward or reject stale requests rather than keep an independent writable copy. After D accepts new writes, routing back to B requires reverse catch-up; retaining B’s old snapshot alone does not make rollback safe.”
Interviewer follow-up
Would a cache need the same durable transfer?
Reveal the follow-up answer
“It may refill from its source of truth instead, but I would protect that origin from a mass cold-cache event.”
What the answer must demonstrate: Do not mistake a new ownership map for a complete migration.
Blank-page exercise · 12 minutes
Build the answer yourself
Draw a ring from 0 to 99 with A20, B50, and C80. Place hashes 12, 35, 45, 65, and 90. Add D40, then remove B50.
- Identify every owner before and after each change.
- Explain why hash 90 wraps to A.
- Separate expected movement at large scale from this exact example.
- Explain how writes and stale routing are handled during migration.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Consistent hashing and virtual nodesWhat moves when D40 joins?Recall first, then reveal
Keys in (20,40] change from B50 to D40; P35 moves, while P45 stays with B.
A new token takes only its preceding interval.
Return to lessonConsistent hashing and virtual nodesWhy does P90 belong to A20?Recall first, then reveal
Clockwise search wraps from 99 to 0 and first reaches A20.
The number line closes into a circle.
Return to lessonConsistent hashing and virtual nodesAre virtual nodes extra copies?Recall first, then reveal
No. Several ring positions can belong to one physical machine; redundant copies need separate placement.
Many tokens are not many hosts.
Return to lessonConsistent hashing and virtual nodesDoes balanced key placement solve a viral photo hot key?Recall first, then reveal
No. A single hot key can dominate requests even when key counts are evenly distributed.
Hash spreads keys, not one key’s popularity.
Return to lessonFinal revision
Summary and interview notes
Consistent hashing keeps most keys on their existing machines when machines join or leave. Virtual nodes give each machine several smaller ranges. This reduces copying or cache refill work. Replication, safe data transfer, full key identity and heavily requested keys still need separate handling.
Remember these points
- A key belongs to the first clockwise token, including an exact match and wraparound.
- Adding D40 between A20 and B50 moves only (20,40]; hash 35 moves, hash 45 stays.
- The expected 1/(N+1) movement on addition assumes suitable balanced placement; a particular insertion can differ.
- Virtual tokens are not physical replicas, and balanced key counts do not guarantee balanced bytes or request rates.
- After routing cutover, rollback must preserve writes accepted by the new owner.
Interview tips
- Compute both an ordinary key and a wraparound key before discussing virtual nodes.
- Separate movement of primary ownership from copying bytes and from changing replica placement.
- Compare a ring with fixed logical buckets when the interviewer asks whether consistent hashing is required.
Important qualifications
- Placement-hash collisions do not merge records; retain and compare complete object keys.
- All routers need compatible hashing and membership versions.
- The six-token example is independent of the original D40 insertion example and deliberately remains imperfectly balanced.
Technical references
- Dynamo: Amazon’s Highly Available Key-value StorePrimary reference for consistent hashing, virtual-node placement, and replica placement choices. Numeric ring examples are original.
- Apache Cassandra: Dynamo ArchitectureOfficial token, virtual-node, and distinct-physical-replica placement description; no prescribed token count or latest-version claim.
- Thaler and Ravishankar: A Name-Based Mapping Scheme for RendezvousOriginal highest-random-weight placement paper. Scores in the chapter are an illustrative calculation, not benchmark output.
- RFC 8584: Highest Random Weight algorithmStandards-track description of object/node scoring, deterministic placement and limited remapping; the chapter uses a general placement example rather than EVPN configuration.
Practice marks stay in this browser.