Skip to content

Explain consistent hashing and virtual nodes

AdvancedAsked very oftenSystem designConceptScalability
#consistent-hashing#sharding#virtual-nodes#scalability

What interviewers are testing

This question separates candidates who memorized "add virtual nodes" from those who can do the distribution math: why modulo hashing remaps 1 minus 1/n of keys, why the ring remaps only about K/n, and how vnodes turn a few arbitrary arcs into a smooth distribution. The follow-ups about failures and rebalancing test whether the model survives contact with operations.

Mental model

Map servers and keys onto the same circular hash space, then assign each key to the first server clockwise. When membership changes, only the arcs that changed hands move — about K/n keys instead of K times (1 minus 1/n) for modulo. Virtual nodes multiply each server's presence on the ring so small clusters average out.

Step-by-step solution

Step 1 of 5

Place keys on the ring

Consistent hashing places both nodes and keys on a circular hash space, usually 0 to 2^32. Each server is hashed to a position on the ring, and each key is assigned to the first node clockwise from the key's own hash. No central directory is needed: any participant can compute ownership from the ring and the same hash function. The catch is variance. Three physical nodes land at arbitrary points, so the arcs between them are wildly uneven — one node can own 52 percent of the keyspace while another owns 17 percent. Small clusters are the worst case because there are too few sample points for the distribution to average out. Watch the animation as server names and keys are projected onto the ring and each key walks clockwise to its owner. That imbalance is exactly what virtual nodes fix later; first let the simple ring expose the problem.

Animation — Place keys on the ring

Client
from Client
Hash Ring

0 to 2^32

Node A

arc 52%

Node B

arc 31%

Node C

arc 17%

1/5

Hash each server name to a point on the 0 to 2^32 ring.

Edge cases & traps

  • Too few virtual nodes leave visible imbalance: use 100-200 vnodes per physical node and review real key counts, not just ring theory.
  • Poor vnode hashing collides or clumps points: use a strong hash such as 128-bit MurmurHash or xxHash and deduplicate positions.
  • Hot keys defeat any hashing scheme because one key cannot be split: cache them, shard them with a random suffix, or replicate them separately.
  • Rebalancing without replication risks data loss when a node dies mid-transfer: keep R copies on the next clockwise nodes and repair asynchronously.
  • Clients with stale ring views briefly route to the wrong node: version the ring through gossip or a config service and retry lookups on miss.

Follow-up questions

Go deeper: Explore the System Design visualizer