System designby Learnastra

Concept lesson · Foundations

Consistent hashing and virtual nodes

By Anup Rai

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.

The visual modelConsistent hashing: ownership before and after adding a node

Keys go clockwise to the first node token. Adding D at 40 moves (20,40] from B to D; other ranges keep their owners.

Consistent hashing: ownership before and after adding a nodeKeys go clockwise to the first node token. Adding D at 40 moves (20,40] from B to D; other ranges keep their owners. 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.Hash space 0...99; walk clockwise to the ownerA 20D 40B 50C 80key 35key 450 / 100Before: 35 belongs to BAfter: 35 belongs to D45 stays with BBlue arc = moved keys(20,40] onlyD is the new nodePlacement is not replication. Virtual nodes spread ranges; a hot key still needs care.
Read the diagram step by step
  1. Tokens are positions on a hash space, not geographic servers.
  2. Initially A=20, B=50 and C=80. Key 35 belongs to B.
  3. Adding D=40 transfers only keys in (20,40] from B to D. Key 45 still belongs to B.
  4. 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

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 practice

Useful foundations: Data partitioning and sharding · Caching: cache hits, misses, write policies and invalidation

Workload and timing examples are interview assumptions.

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.

Concept in focusAdd one owner: watch key 55 move

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.

Add one owner: watch key 55 moveNode 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. 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.Before: three ownersA20B50C80key 55clockwiseAfter: add D at 60A20B50C80D60key 55clockwise55 goes to C8055 now goes to D60Only keys in (50, 60] change owner: C gives that interval to D. A and Bkeep their ranges.

Remember: Only the new owner’s predecessor interval moves.

Read the diagram
  1. Before the change, A20, B50 and C80 own the ring. Key 55 reaches C80 clockwise.
  2. After D60 joins, key 55 reaches D60 first.
  3. 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.

Worked example diagramFive keys on a numerically scaled hash ring

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.

Five keys on a numerically scaled hash ringA20, 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. 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.Hash space 0-99: follow the ring clockwise to the first owner0 / 100Large dots: node positionsSmall gold dots: keysA20D40 (new)B50C80P12P35P45P65P90Key ownershipKeyBefore D40After D40P12A20A20P35B50D40P45B50B50P65C80C80P90A20A20Only P35 changes owner.The other four keys stay put.Blue arc: (20,40] moves from B to D. D owns 40; A still owns 20.
Read the key assignments
  1. P12 hashes to 12: A20 before adding D40; A20 afterward.
  2. P35 hashes to 35: B50 before adding D40; D40 afterward.
  3. P45 hashes to 45: B50 before adding D40; B50 afterward.
  4. P65 hashes to 65: C80 before adding D40; C80 afterward.
  5. 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.

Concept in focusSix virtual nodes, three physical servers

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.

Six virtual nodes, three physical serversA1 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. 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.Hash positions 0-99Small black dots: example keysPhysical servers0 / 100A1 at 10B1 at 30C1 at 45A2 at 60B2 at 80C2 at 95clockwisekey 5key 55Physical server ATokens: A1 (10), A2 (60)(95,10] and (45,60]30 positions = 30%Physical server BTokens: B1 (30), B2 (80)(10,30] and (60,80]40 positions = 40%Physical server CTokens: C1 (45), C2 (95)(30,45] and (80,95]30 positions = 30%Each colored arc ends at the token that owns it.Two lookups, one physical destinationHash 5A1 at 10Hash 55A2 at 60Physical server Astores both keysA1 and A2 share server A's storage and failure risk.Replication requires other physical owners.

Remember: Several ring positions can point to one physical server.

Read the diagram
  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. 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.
  6. 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.

Foundation · Question 1

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.

What the answer must demonstrate: Demonstrate the rule with actual positions.

Applied · Question 2

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.”

What the answer must demonstrate: Keep a concrete trace distinct from a statistical estimate.

Applied · Question 3

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.”

What the answer must demonstrate: Placement and durability are separate responsibilities.

Foundation · Question 4

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.”

What the answer must demonstrate: Avoid claiming every modulo mapping necessarily changes.

Foundation · Question 5

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.”

What the answer must demonstrate: Count physical failure domains for replication.

Applied · Question 6

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.”

What the answer must demonstrate: Key count, bytes, and traffic are different load measures.

Follow-up · Question 7

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.”

What the answer must demonstrate: Qualify both arithmetic and assumptions.

Follow-up · Question 8

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.”

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 lesson
Consistent 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 lesson
Consistent 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 lesson
Consistent 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 lesson

Final 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

Practice marks stay in this browser.