System design · Rebalancing with minimal movement
Consistent hashing explained
Any system that spreads data across machines needs a rule for which machine owns which key. The obvious rule - hash the key and take it modulo the number of servers - works until you add or remove a server. Then almost every key moves. Consistent hashing is the fix: when the cluster changes size, only a small, predictable share of keys moves.
Updated · 6 min read
The problem with hash mod N
With N servers, the simple rule is server = hash(key) % N. It spreads keys evenly and needs no lookup table. The trouble starts when N changes. A key stays put only if hash % N and hash % (N+1) give the same answer, and for most keys they don't. Going from N to N+1 servers moves about N/(N+1) of all keys. Consistent hashing moves about 1/(N+1) - only the keys the new server takes over.
| Change | Keys moved with hash mod N | Keys moved with a hash ring |
|---|---|---|
| 4 → 5 nodes | 4/5 = 80% | 1/5 = 20% |
| 10 → 11 nodes | 10/11 ≈ 91% | 1/11 ≈ 9% |
| 100 → 101 nodes | 100/101 ≈ 99% | 1/101 ≈ 1% |
| 5 → 4 nodes (one fails) | 4/5 = 80% | 1/5 = 20% - only the failed node's keys |
For a cache, moving 80% of keys means an 80% miss rate right after a deploy or a crash, and that load lands on the database at once. For a database, it means copying most of the data across the network.
The hash ring
Picture the output range of a hash function - say 0 to 2^32 - 1 - bent into a circle. Hash each server's ID to a point on that circle. To find a key's owner, hash the key to a point and walk clockwise until you hit a server. That server owns the key.
Now add a server. It lands at one point and takes over only the keys between it and the previous server counter-clockwise. Every other key keeps its owner. Remove a server and its keys pass to the next server clockwise; nothing else moves.
In code, the ring is a sorted array of positions, and a lookup is a binary search that wraps to the start.
ring = sorted list of (position, node)
function lookup(key):
h = hash(key)
lo, hi = 0, len(ring)
while lo < hi:
mid = (lo + hi) / 2
if ring[mid].position < h:
lo = mid + 1
else:
hi = mid
if lo == len(ring):
lo = 0
return ring[lo].nodeLookup is O(log M), where M is the number of points on the ring. Clients can hold the ring locally, so routing needs no network call - only a fresh copy of the ring when membership changes.
Virtual nodes
With one point per server, the arcs between servers are random lengths. One server can easily own several times as many keys as another. And when a server fails, its whole arc goes to a single neighbour, which may then fall over too.
Virtual nodes fix both problems. Each physical server gets many points on the ring - for example by hashing "node-a#0", "node-a#1" and so on. With 100 to 200 points per server, each server owns many small arcs scattered around the ring, and the totals even out.
- Even load: many small arcs average out, so each server's share stays close to 1/N.
- Spread recovery: a failed server's arcs are scattered, so its keys spread across many survivors instead of one.
- Different machine sizes: give a server with twice the memory twice the virtual nodes, and it gets about twice the keys.
- The cost: a bigger ring to store and search, and more small ranges to track when moving data. Cassandra lowered its default token count per node in version 4.0 for this reason.
Replication on the ring
The ring also tells you where copies go. To store R replicas of a key, start at the key's position and walk clockwise, taking the first R distinct physical servers you meet. With virtual nodes, skip any point that belongs to a server you have already picked, or two replicas could land on the same machine.
This list of R servers is what Dynamo calls the preference list. Any of them can serve the key, and when one fails, the next server along the ring takes its place.
Where it's used
- Dynamo-style key-value stores: Amazon's Dynamo paper made the ring, virtual nodes and preference lists standard.
- Cassandra: each node owns token ranges on a ring, and replicas are placed by walking the ring with rack and datacenter awareness.
- Memcached client libraries: the servers know nothing about each other. The client hashes keys to servers with a ring - the ketama scheme - so adding a cache server doesn't wipe the cache.
- Load balancers: consistent-hash routing sends the same user or session to the same backend. Nginx, HAProxy and Envoy all support it.
Alternatives
- Rendezvous hashing (highest random weight): for each key, score every server with hash(key, server) and pick the highest. Adding or removing a server moves only about 1/N of keys, it needs no virtual nodes, and the top R scores give you a replica list. The cost is O(N) work per lookup, which is fine for tens of servers.
- Jump consistent hash: a short function from Google that maps a key to one of N numbered buckets with no ring and no memory, and moves only 1/(N+1) of keys when a bucket is added. The catch is that buckets can only be added or removed at the end, so it suits numbered shards more than servers that fail at random.
What interviewers look for
- Explaining why hash mod N fails, with numbers - most keys move when N changes.
- Drawing the ring and showing that a change moves only about 1/N of keys.
- Bringing up virtual nodes unprompted, for even load and for mixed machine sizes.
- Placing replicas by walking the ring to distinct physical servers.
- Knowing its limit: consistent hashing balances keys, not traffic. A hot key - one celebrity, one viral product - still lands on one server. That needs separate handling: a local cache in front, replicating the hot key to more servers, or splitting it into sub-keys like key#1 to key#10 and spreading reads across them.
Frequently asked questions
What is consistent hashing in simple terms?
+
A way to assign keys to servers so that adding or removing a server moves only a small share of keys. Servers and keys are hashed onto the same circle, and each key belongs to the first server clockwise from it.
How many keys move when you add a node?
+
With consistent hashing, about 1/(N+1) of keys when going from N to N+1 nodes - only the keys the new node takes over. With hash mod N, about N/(N+1) move, which is 80% when going from 4 to 5 nodes.
Why do we need virtual nodes?
+
With one point per server, arcs on the ring have uneven lengths, so load is uneven, and a failed server dumps all its keys on one neighbour. Many points per server even out the load, spread a failed server's keys across the cluster, and let bigger machines take a bigger share.
Does consistent hashing solve hot keys?
+
No. It spreads keys evenly, but a single very popular key still maps to one server. Hot keys need their own fix, such as caching them closer to clients, replicating them more widely, or splitting them into several sub-keys.
Consistent hashing or rendezvous hashing?
+
Both move only about 1/N of keys on a change. Rendezvous hashing is simpler and needs no virtual nodes, but each lookup checks every server, so it fits small clusters. A ring with binary search scales better to large clusters with many lookups.