Consistent Hashing Explained (with Diagrams)
Consistent hashing is one of a handful of ideas that shows up in almost every distributed system you will ever be asked to design — caches, sharded databases, load balancers, object stores. It is also one of the most commonly botched interview answers, because most candidates can recite "keys go on a ring" without being able to say what problem the ring actually solves or why the naive approach fails.
This post builds it from the problem up: what breaks, why the ring fixes it, the arithmetic behind the "only 1/N keys move" claim, why plain consistent hashing is unusable in production without virtual nodes, and what it still does not solve.
Start with the thing that breaks
You have a distributed cache — four servers, and you need to decide which server holds which key. The obvious approach is modulo hashing:
server_index = hash(key) % N # N = number of servers
With N = 4 this is clean, fast, and distributes keys evenly. It also falls apart the first time your cluster changes size.
Take six keys with these hash values and watch what happens when you add a fifth server:
You can check the arithmetic yourself: 17 % 4 = 1 but 17 % 5 = 2. 56 % 4 = 0 but 56 % 5 = 1. Only 23 and 42 happen to land on the same index both times.
That "roughly one in five" is not a coincidence. For a random key, hash(key) % 4 and hash(key) % 5 are effectively independent, so the chance they agree is about 1/5. Generalising: moving from N to N+1 servers leaves about 1/(N+1) of keys in place and moves all the rest.
For a cache, this is a catastrophe rather than an inconvenience. Every relocated key is a miss. Every miss goes to the database. Adding one cache server to relieve load causes roughly 80% of your traffic to stampede the database simultaneously — so the capacity increase you deployed to survive a traffic spike is itself the outage. The same thing happens in reverse when a server dies and N drops to 3.
The flaw is structural: modulo hashing binds every key's location to the exact server count. In a system where nodes are added, removed, and fail routinely, that binding has to go.
The ring
Consistent hashing breaks the dependency on N by giving servers and keys positions in the same space, rather than computing an index from a count.
- Take the output range of your hash function — say 0 to 2³²−1 — and treat it as a circle, so the largest value sits next to 0.
- Hash each server (by hostname, IP, or an assigned ID) to a point on that circle.
- Hash each key to a point on the same circle.
- To find a key's owner, walk clockwise from the key's position until you hit a server. That server owns the key.
Notice what is not in that procedure: the number of servers. A key's position never changes. Only the set of positions it might walk to does.
Why adding and removing nodes is now cheap
This is the property worth being able to state precisely in an interview, because it is the entire payoff.
Remove a server. If S2 dies, the keys that used to stop at S2 now continue clockwise to the next server along. Every other key on the ring is untouched — its clockwise walk was never passing through S2. You remap the keys in one arc, about 1/N of the total.
Add a server. Insert S4 somewhere on the ring and it claims exactly the keys lying between it and the previous server counter-clockwise. Those keys used to belong to whichever server is now immediately clockwise of S4. Nobody else is affected.
The headline number: a cluster change moves about 1/N of keys rather than about (N−1)/N. Going from 4 servers to 5, that is 20% versus 80%. At 20 servers it is 5% versus 95%. The bigger your cluster, the more dramatic the difference — which is exactly backwards from how modulo hashing behaves.
For a cache, this is the difference between a brief dip in hit rate and a database outage.
The catch nobody mentions: the ring is lumpy
Here is where most interview answers stop, and where a good follow-up question lives. Plain consistent hashing, exactly as described so far, does not distribute load evenly, and the reason is uncomfortable once you see it.
Three servers hashed to random points on a circle do not divide it into three equal thirds. They divide it into three arcs of random size. On average each owns a third, but the variance is enormous: with N servers placed at random, the largest arc is typically several times the size of the smallest. One server gets hammered while another idles.
It gets worse under failure. When a server dies, its entire arc transfers to exactly one neighbour — the next server clockwise. That neighbour's load can double instantly. If it was already near capacity, it falls over too, and its arc lands on its neighbour. That is a cascading failure with a very short fuse.
So plain consistent hashing solves the remapping problem and introduces a load-distribution problem.
Virtual nodes fix both
The fix is almost embarrassingly simple: stop placing each server on the ring once. Place it many times.
For each physical server, hash it together with a replica index — hash("cache-03#0"), hash("cache-03#1"), … hash("cache-03#159") — and put a point on the ring for each. These are virtual nodes (Dynamo calls them tokens). A key still walks clockwise to the nearest point; that point maps back to its physical owner.
Virtual nodes buy three things at once:
- Even distribution. Many small arcs per server average out. The relative standard deviation of load falls roughly as 1/√V, where V is the number of virtual nodes per server.
- Graceful failure. A dead server's arcs are scattered around the ring, so its load is absorbed by every remaining server in proportion, not dumped on one.
- Heterogeneous hardware. Give a machine with twice the memory twice the virtual nodes and it takes twice the keyspace. Plain consistent hashing has no way to express that at all.
Roughly what V buys you:
| Virtual nodes per server | Approx. load spread | Ring size at 10 servers |
|---|---|---|
| 1 | ±100% or worse | 10 points |
| 10 | ~±32% | 100 points |
| 100 | ~±10% | 1,000 points |
| 500 | ~±4.5% | 5,000 points |
The cost is memory and lookup time — the ring is a sorted structure you binary-search, so it grows with N × V. A few thousand entries is nothing; a few million starts to matter. libketama, the implementation behind most Memcached clients, settled on 160 points per server. Amazon's Dynamo paper describes 100–200 tokens per node. Cassandra's num_tokens historically defaulted to 256 and now defaults to 16, paired with a smarter allocation algorithm that places tokens deliberately rather than randomly.
That last detail is worth carrying into an interview: the modern trend is fewer, better-chosen virtual nodes rather than more random ones, because random placement needs large V to behave, and large V makes operations like bootstrapping a new node slower.
What the implementation actually looks like
Small enough to write on a whiteboard, which is why interviewers ask for it:
import bisect
from hashlib import blake2b
def h(s: str) -> int:
return int.from_bytes(blake2b(s.encode(), digest_size=8).digest(), "big")
class Ring:
def __init__(self, replicas: int = 160):
self.replicas = replicas
self.points: list[int] = [] # sorted positions
self.owner: dict[int, str] = {} # position -> physical server
def add(self, server: str) -> None:
for i in range(self.replicas):
p = h(f"{server}#{i}")
bisect.insort(self.points, p)
self.owner[p] = server
def remove(self, server: str) -> None:
for i in range(self.replicas):
p = h(f"{server}#{i}")
self.points.remove(p)
del self.owner[p]
def get(self, key: str) -> str:
if not self.points:
raise LookupError("empty ring")
i = bisect.bisect(self.points, h(key))
return self.owner[self.points[i % len(self.points)]] # wrap at the end
Three details carry all the meaning. bisect gives O(log NV) lookup. The % len(self.points) is the wrap — a key past the last point belongs to the first. And the ring is derived: every client that hashes the same server names builds an identical ring independently, with no coordination. That last property is why consistent hashing works for client-side sharding, where there is no router to ask.
What it does not solve
Being able to name the limits is what separates a memorised answer from an understood one.
Hot keys. Consistent hashing distributes keys evenly, not traffic. If one key is red-hot — a celebrity's profile, a link that just went viral — every request for it still lands on one server. The ring cannot help; you need replication of that key, a small local cache in front, or request coalescing. This is a very common follow-up, and it catches people who have only memorised the ring.
Replication and consistency. The ring tells you where a key lives. It says nothing about how many copies exist or how they stay in sync. Dynamo-style systems layer replication on top by walking clockwise past the first node and storing to the next R distinct physical servers.
The cost of moving data. "Only 1/N of keys move" is cheap for a cache, where a moved key is simply a miss. For a database it means actually streaming a shard's worth of data between machines while serving traffic. The remapping is minimal; the migration is not free.
Bounded load under skew. Even with virtual nodes, load is only statistically even. Google's "consistent hashing with bounded loads" adds a hard cap: a key overflows to the next node once its intended owner exceeds a multiple of the average. HAProxy and Envoy both implement it.
The alternatives worth naming
Mentioning one of these unprompted is a strong signal, provided you can say when you would pick it.
Rendezvous hashing (highest random weight, 1997). For each key, compute hash(key, node) for every node and pick the highest. No ring, no virtual nodes, and a naturally even distribution. The cost is O(N) per lookup instead of O(log N) — fine for tens of nodes, not for thousands. It is simpler to implement correctly and handles weighting cleanly.
Jump consistent hash (Google, 2014). Seven lines, no memory at all, perfectly even distribution. The catch is severe: it maps to buckets numbered 0..N-1, so you can only add or remove at the end of the range. Losing an arbitrary node in the middle is not expressible. Great for sharding into a fixed numbered set, useless for a cluster where any machine can die.
Maglev hashing (Google, 2016). Builds a fixed-size lookup table so lookups are a single array index, trading a small amount of extra disruption on node change for speed and very even distribution. Designed for load balancers doing this per-packet.
Using it in an interview
Reach for consistent hashing when you are partitioning data or load across a set of nodes that changes — distributed caches, sharded key-value stores, stateful services that need sticky routing without a shared session store. If the node set is genuinely fixed, say so and use something simpler; proposing a ring for a three-node cluster that never changes is a sign you are pattern-matching rather than designing.
A complete answer, in the order that lands best:
- Name the failure first. "Modulo hashing remaps almost every key when the server count changes — for a cache that's a full-miss storm against the database."
- Describe the ring in one breath. Servers and keys hashed into the same space, each key owned by the first server clockwise.
- Quantify the win. About 1/N of keys move on a cluster change, and only one neighbour is involved.
- Volunteer the weakness. Random placement gives lumpy arcs and dumps a dead node's entire load on one neighbour.
- Close it with virtual nodes, and say what V you would use and why.
- Pre-empt the hot-key follow-up. Say that the ring balances keys, not traffic, and name what you would do about it.
Steps 4 and 5 are the ones most candidates skip, and they are the ones that make the difference — anyone can describe a ring, but volunteering the flaw in your own answer before the interviewer finds it is exactly the behaviour senior interviewers are scoring for.
If you want to pressure-test how this sounds out loud — including the follow-ups about hot keys and replication that always come next — you can run a full system design mock by voice on Whitepad and get scored on it.
Practice this out loud
Reading is the easy part. Sit across from a senior AI interviewer that talks, watches your whiteboard, and scores you like the real thing — your first mock is free.
Start a free mock →