← All posts
ConceptsAugust 15, 202610 min read

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:

hash(key) 17 23 42 56 71 88 N = 4 before S1 S3 S2 S0 S3 S0 N = 5 after S2 S3 S2 S1 S1 S3 moved stayed stayed moved moved moved
One extra server, four of six keys relocated. Nothing about the keys changed — only the modulus did. Across a full keyspace, going from four servers to five leaves roughly one key in five where it was.

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.

  1. 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.
  2. Hash each server (by hostname, IP, or an assigned ID) to a point on that circle.
  3. Hash each key to a point on the same circle.
  4. To find a key's owner, walk clockwise from the key's position until you hit a server. That server owns the key.
0 / 2³² S1 S2 S3 key A key B key C clockwise key A S2 key B S3 key C S1 Key C sits just past S3, so its clockwise walk crosses the 0 point and lands on S1. The ring wraps — there is no first or last server.
Position replaces arithmetic. A key's owner is decided by where it falls relative to the servers, not by how many servers exist. That single change is what makes the cluster resizable.

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.

S1 S2 S3 S4 key A this arc transfers from S2 to S4 Key A moves. It used to walk past this point to reach S2; now S4 stops it. Every other key stays. S1 and S3 are not involved at all — no coordination, no rebalance, no cache-wide invalidation.
One neighbour is affected, not the whole cluster. Compare this with the first figure, where a single added server rewrote four of six assignments. Here the blast radius is a single arc, and it is bounded by construction.

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.

cache-01 cache-02 cache-03 cache-04 Four servers, six positions each. Each server now owns six small arcs instead of one large one, so the totals even out. When cache-03 dies, its six arcs go to six different neighbours — roughly a third each, not double for one. Real deployments use 100–200 points, not 6.
The same four servers, interleaved. Six points each is enough to show the mechanism; production uses one to two hundred. The drawing is what makes the failure behaviour obvious — a dead server's share is scattered, so no single neighbour inherits it.

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 serverApprox. load spreadRing size at 10 servers
1±100% or worse10 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:

  1. 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."
  2. Describe the ring in one breath. Servers and keys hashed into the same space, each key owned by the first server clockwise.
  3. Quantify the win. About 1/N of keys move on a cluster change, and only one neighbour is involved.
  4. Volunteer the weakness. Random placement gives lumpy arcs and dumps a dead node's entire load on one neighbour.
  5. Close it with virtual nodes, and say what V you would use and why.
  6. 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 →