Skip to content
BytePatterns

Consistent Hashing Explained: Why Keys Barely Move

8 min readBytePatterns

Adding a fifth cache server with key modulo n remaps 80% of keys; on a hash ring it moves about 20%. The ring, virtual nodes and the failure case, measured.

Every distributed cache, sharded key-value store and partitioned queue has to answer the same small question for every request: which machine owns this key? The obvious answer works perfectly until the day the number of machines changes — and then it quietly turns a routine scaling event into a stampede against the database behind it.

Consistent hashing is the fix, and it is worth being able to derive rather than just name.

The problem it solves

Four cache servers. The natural rule is hash(key) % 4: cheap, stateless, perfectly even. Every client computes the same owner without asking anyone.

Now traffic grows and you add a fifth server. The rule becomes hash(key) % 5, and for most keys the answer changes. A key whose hash is 12 lived on server 0 and now lives on server 2. Most keys are suddenly looked up on a machine that has never seen them. The hit rate collapses at the exact moment you added capacity because load was high, and every miss falls through to the origin.

The same happens in reverse when a server dies: % 4 becomes % 3, and the survivors lose most of their useful contents too.

The intuition

The modulo rule ties every key to the count of machines, so changing the count reshuffles everything. The fix is to stop doing that.

Picture the hash space as a circle. Hash each server onto it, at some point. Hash each key onto the same circle. A key is owned by the first server you meet walking clockwise from the key's point.

Now add a server. It lands at one point and takes over only the arc between itself and the previous server counter-clockwise. Keys on that arc move to it; every other key keeps exactly the owner it had. Remove a server and the reverse happens: its arc falls to the next server clockwise, and nobody else is touched.

On a ring, a change to one server moves one server's share of keys. Under modulo, a change to one server moves nearly all of them.

There is a catch. With one point per server, the arcs are whatever lengths the hash happens to produce, and they are rarely similar. One server can own half the ring. The standard fix is virtual nodes: hash each server onto the ring many times — n0#0, n0#1, … — so it owns many small arcs instead of one big one. Their lengths average out, and when a server dies its arcs are scattered, so its load is spread over all the survivors instead of landing on a single neighbour.

Watch it run

The animation follows one key onto the ring, to its owner, through a miss. Then it compares the two rules when a fifth box arrives — the readout goes from "almost all" remapped under modulo to "one share" on the ring — then a machine dying, virtual points evening out the arcs, and eviction.

Design a Distributed Cache

Step 1 of 11

One working set, too large for a single box, spread over machines arranged as a ring.

The same interactive animation as the lesson — step through it with the controls.

The code

A ring is a sorted list of points plus a binary search. bisect finds the first server point clockwise of the key's hash; running off the end wraps back to the first point.

import bisect, hashlib

def h(s):                                     # stable across runs, unlike hash()
    return int(hashlib.md5(s.encode()).hexdigest()[:8], 16)

class Ring:
    def __init__(self, nodes, vnodes=1):
        self.points = sorted((h(f"{n}#{v}"), n) for n in nodes for v in range(vnodes))
        self.hashes = [p for p, _ in self.points]
    def owner(self, key):
        i = bisect.bisect(self.hashes, h(key))       # first point clockwise
        return self.points[i % len(self.points)][1]  # past the end: wrap around

keys = [f"user:{i}" for i in range(100_000)]

def moved(before, after):
    return sum(before(k) != after(k) for k in keys) / len(keys)

mod4 = lambda k: f"n{h(k) % 4}"
mod5 = lambda k: f"n{h(k) % 5}"
print(f"{moved(mod4, mod5):.0%}")                          # 80%

old = Ring(["n0", "n1", "n2", "n3"], vnodes=100)
new = Ring(["n0", "n1", "n2", "n3", "n4"], vnodes=100)
print(f"{moved(old.owner, new.owner):.0%}")                # 21%
print(all(old.owner(k) == new.owner(k)                     # every mover went to n4
          for k in keys if new.owner(k) != "n4"))           # True

less = Ring(["n0", "n1", "n3"], vnodes=100)                # n2 dies
print(all(less.owner(k) == old.owner(k)                    # only n2's keys move
          for k in keys if old.owner(k) != "n2"))           # True

Going from four to five servers under modulo moved 80% of 100,000 keys. On the ring it moved 21% — close to the ideal fifth, since a new server should end up with about a fifth of the keys — and every key that moved went to the new server. No key was shuffled between two old ones. Removing n2 touched only the keys n2 had owned.

Python's built-in hash() is randomised per process for strings, which is why the code uses a stable hash: every client must agree on where a key lands.

Virtual nodes, measured. With one point per server the four loads are wildly uneven, and a dying server hands everything to one neighbour. With 100 points each, both problems go away:

from collections import Counter
for v in (1, 100):
    before = Ring(["n0", "n1", "n2", "n3"], vnodes=v)
    after = Ring(["n0", "n1", "n3"], vnodes=v)
    load = Counter(map(before.owner, keys))
    heirs = Counter(after.owner(k) for k in keys if before.owner(k) == "n2")
    print(v, sorted(load.values()), sorted(heirs.items()))
# 1 [9866, 17431, 22569, 50134] [('n3', 9866)]
# 100 [24060, 25157, 25188, 25595] [('n0', 9307), ('n1', 6933), ('n3', 7820)]

def brute_owner(ring, key):                   # scan every point, no bisect
    after = [pt for pt in ring.points if pt[0] > h(key)]
    return min(after or ring.points)[1]
print(all(brute_owner(r, k) == r.owner(k)
          for r in (old, new, less) for k in keys[::50]))   # True

One point each: one server holds half of all keys, another under a tenth. A hundred points each: every server within a few percent of 25,000. The last check compares the bisect lookup, wrap-around included, with a plain scan over every point.

The complexity

  • Lookup: O(log(N · V)) for N servers with V virtual nodes each — one binary search over the sorted points.
  • Adding or removing a server: insert or delete V points; about 1/N of the keys change owner, which is the minimum possible, since the new server has to end up with its share from somewhere.
  • Memory: N · V points, held by every client. A few hundred points per server is typical and trivially small.

Where it goes wrong

  • Too few virtual nodes. One point per server gave a 5:1 load imbalance above. Uneven arcs mean one hot server, which is the bottleneck you added servers to avoid.
  • An unstable hash. Clients that disagree on a key's hash disagree on its owner. Use a fixed hash function, never a per-process randomised one.
  • Forgetting that moved keys are cold. Consistent hashing limits how many keys move; it does not warm them. A new server still starts empty, so add capacity before you need it, not during the peak.
  • Heterogeneous machines. A server with twice the memory should get twice the virtual nodes. Equal V assumes equal machines.
  • Replication. Stores that keep copies usually place them on the next distinct servers clockwise; skip further points that belong to a server already chosen, or two replicas can land on one machine.

The replica placement and quorum side of this is its own design: see designing a key-value store, and database sharding for the partitioning trade-offs around it.

How to say it in an interview

"hash mod N remaps almost every key when N changes, so a scaling event empties the cache. Instead I'd put servers on a hash ring and send each key to the next server clockwise. Adding or removing a server then moves only about 1/N of the keys — the ones on its arc. To keep load even and to spread a failed server's keys across the survivors, each server gets a few hundred virtual nodes. Lookup is a binary search over the sorted points."

Then add the caveat: the moved keys still start cold, so the ring limits the damage of a change — it does not make it free.