Day 6 B-Building blocks 2026-10-01 ← All lessons

Consistent hashing: why hash % N breaks and how to fix it

Add one cache server under hash % N and most of your keys now live somewhere else. Put the servers on a ring and only one neighbour's slice moves.

160SHA-1 hash space used for the ring: 0 to 2^160 - 1
10Load standard deviation with 100 virtual nodes per server (of the mean)
5Load standard deviation with 200 virtual nodes per server (of the mean)

Key points

Outcomes

01Explain out loud why modulo placement falls apart when the pool grows or shrinks, with real remap counts.
02Walk through a lookup, a join and a leave on the ring, and name exactly which keys move and where they go.
03Argue the virtual node tradeoff, and say which problems it does not solve, hot keys in particular.
01

The problem: hash % N is a promise about N


Say you run 4 cache servers. The obvious way to pick one is serverIndex = hash(key) % N: hash the key and take the remainder. The book's example uses 4 servers and 8 keys. hash(key0) % 4 = 1, so key0 lives on server 1, and the keys spread out nicely. As long as N stays fixed and the hash is uniform, this works.

Now server 1 dies and N becomes 3. The hashes have not changed, but you are dividing by a different number, so nearly every remainder changes. In the book's example most keys get a new server, not only the ones that lived on server 1. Every client now asks the wrong server, gets a miss, and falls through to the database. That is the cache miss storm. Losing one box costs you far more than one box of cache, and the extra load lands on the database right when you have lost capacity.

Here are real numbers. The script in the code section hashes 10,000 keys and grows the pool from 4 servers to 5. Modulo placement moves 7,975 of them, about 80%. A consistent hashing ring with the same keys moves 1,985, about 20%. That is roughly the one fifth the new server should own, and nothing more.

hash % N
7,975 moved
Consistent hashing ring
1,985 moved
Keys remapped out of 10,000 when going from 4 to 5 servers (computed by the script in the code section).
GotchaGotcha: failures are not the only time you pay. Scaling up for a traffic spike also changes N, so the deploy meant to add capacity flushes your cache at the moment you need it most.
02

The mental model: a clock face with servers on it


Take the full output range of a hash function and bend it into a circle. With SHA-1 that range runs from 0 to 2^160 - 1, and the two ends meet. Hash each server's name or IP onto the circle. Hash each key onto the same circle with the same function. There is no modulo anywhere, so where something sits on the ring never depends on how many servers you have.

The rule: a key belongs to the first server you meet walking clockwise from it. Picture a clock with a few people standing at different hours. A letter dropped at 2 o'clock goes to the next person clockwise. If someone new steps in at 3 o'clock, only the letters between the previous person and 3 o'clock change hands. Every other letter stays put.

  1. Hash every server onto the ring once, when it joins.

  2. Hash the key with the same function.

  3. Walk clockwise from the key's position. The first server you reach owns the key.

  4. If you pass the top of the ring without meeting a server, wrap around to the first one. The ring has no end.

That gives you the textbook guarantee. When the table is resized, only k/n keys need to be remapped on average, where k is the number of keys and n is the number of slots. Modulo moves nearly all k.

03

The mechanism live: joins and leaves


How to read: Press Add a server to bring cache-d onto the ring and watch the moved counter, then compare it with the hash % N meter. Hover any key dot to trace its lookup walk. New keys deals a fresh set of 24 onto the same ring.

Keys that moved on that change
0 of 24

Press a button to make a change.

Same change with hash % N
—

How many keys plain modulo hashing would have relocated.

Hover or tap a key dot on the ring to trace its lookup.
Three cache servers own twenty-four keys. Adding or removing a server moves only the keys in one arc of the ring, and the two meters prove it against plain modulo hashing.

Join. At the start the three servers split 24 keys as 12, 3 and 9. Already lopsided, which is the weakness section below, but watch what happens on change. Press Add a server: cache-d joins and the counter reads 4 of 24. All four, user:3310, user:9025, cart:71ba and post:1160, come from cache-c, the neighbour whose arc cache-d landed inside. cache-a and cache-b are untouched, so their clients keep hitting warm caches. The second meter shows what plain hash % N would have done on the identical change: 19 of 24 keys relocated.

Leave. Remove a server and its whole arc goes to the next server clockwise, and only there. Take cache-a out of the original three and all 12 of its keys land on cache-b, which jumps from 3 keys to 15 in one step. Nothing else moves. That instant doubling is the cascading-failure risk the next section is about: if cache-b was already close to capacity, it falls over next.

  1. Keys affected by a join: start at the new server and walk anticlockwise until you reach the previous server. Keys in that range move to the new server.

  2. Keys affected by a leave: start at the removed server and walk anticlockwise to the previous server. Keys in that range move to the next server clockwise.

  3. Nothing outside that range moves. So rebalancing is a range copy between two neighbours, not a reshuffle of the whole cluster.

Interview tipInterview line: "On a join, only the keys between the new node and the node before it (anticlockwise) move, and they all come from one existing node. On a leave, the dead node's range goes to the next node clockwise."
04

The weakness: one point per server lies about balance


The basic ring has two real problems, and the ring above already shows both.

Uneven partitions. A partition is the stretch of ring between a server and the one before it. With one hashed point per server, these stretches are whatever size the hash happens to give, and every join or leave changes them. The book's example: remove s1 and s2's partition becomes twice the size of s0's and s3's. In the ring above, cache-d only took keys from cache-c. A new server relieves one neighbour, not the whole cluster.

Non-uniform key distribution. If the server points happen to bunch together, one server can own most of the ring. The book shows a layout where server 2 holds most of the keys and servers 1 and 3 hold none. The ring above shows the live version: the opening split is 12, 3 and 9, and if you remove cache-a, cache-b jumps from 3 keys to 15 in one step while cache-c sits idle. If cache-b was already close to capacity, it falls over next, its keys pile onto the following server, and you have a cascade.

GotchaGotcha: consistent hashing spreads keys, not traffic. A single very hot key, like a celebrity's profile, still lands on exactly one server no matter what ring you use. Even placement helps when several hot keys would otherwise collide on one shard. One hot key needs replication or a local cache in front of it.
05

The fix: virtual nodes, and what they cost


Don't hash each server once. Hash it many times under different labels: server 0 becomes s0_0, s0_1, s0_2 and so on, and each label gets its own point on the ring. The book uses 3 per server to illustrate. Real systems use far more. Lookup is unchanged: walk clockwise to the first virtual node, then map it back to its real server.

This works because each server now owns many small arcs spread around the ring, not one big arc of random size, and the randomness averages out. The book cites an experiment where the load standard deviation is about 10% of the mean with 100 virtual nodes and about 5% with 200. It also fixes the failure case. When a server dies, each of its small arcs goes to a different successor, so its load spreads across the cluster instead of landing on one neighbour. A new server takes a little from everyone.

Option A

More virtual nodes

  • Lower load variance: about 10% at 100, 5% at 200
  • A dead server's load spreads over many survivors
  • A new server takes a little from every server
  • Bigger machines can get more vnodes for a larger share
Option B

Fewer virtual nodes

  • Less metadata: every router stores every point
  • Faster ring lookups and smaller membership updates
  • Fewer, larger ranges to move during rebalance
  • Simpler to reason about and debug

The cost is space. Every client or router has to hold the full list of virtual node positions, and every membership change rewrites more entries. The book treats this as a dial: pick the number of virtual nodes that brings variance under your target, and stop there.

06

The code: a ring with virtual nodes


python
import bisect, hashlib

def h(s):
    return int(hashlib.md5(s.encode()).hexdigest()[:8], 16)

class Ring:
    def __init__(self, vnodes=100):
        self.vnodes = vnodes
        self.points = []      # sorted ring positions
        self.owner = {}       # position -> real server

    def add(self, server):
        for i in range(self.vnodes):
            p = h(f"{server}#{i}")
            bisect.insort(self.points, p)
            self.owner[p] = server

    def remove(self, server):
        for i in range(self.vnodes):
            p = h(f"{server}#{i}")
            self.points.remove(p)
            del self.owner[p]

    def lookup(self, key):
        i = bisect.bisect_right(self.points, h(key))
        if i == len(self.points):
            i = 0             # wrap around the ring
        return self.owner[self.points[i]]

keys = [f"user:{n}" for n in range(10000)]
mod_before = {k: h(k) % 4 for k in keys}
mod_after = {k: h(k) % 5 for k in keys}
print("hash % N, 4 -> 5 servers, moved:", sum(mod_before[k] != mod_after[k] for k in keys))

ring = Ring()
for s in ["s0", "s1", "s2", "s3"]:
    ring.add(s)
before = {k: ring.lookup(k) for k in keys}
ring.add("s4")
after = {k: ring.lookup(k) for k in keys}
print("ring, 4 -> 5 servers, moved:", sum(before[k] != after[k] for k in keys))
print("all moved keys went to s4:", all(after[k] == "s4" for k in keys if before[k] != after[k]))

# Output:
# hash % N, 4 -> 5 servers, moved: 7975
# ring, 4 -> 5 servers, moved: 1985
# all moved keys went to s4: True

The lines that matter. h(f"{server}#{i}") is the virtual node trick: one server becomes 100 labels and so 100 points. bisect_right on the sorted list is the clockwise walk, done as an O(log n) binary search, so lookups stay fast however many points you add. The i == len(self.points) check is the wrap past the top of the ring. Leave it out and any key past the last point crashes the lookup. The last print checks the property you are paying for: every key that moved went to the new server, and none moved between old servers.

Production notes: real systems keep this table in a membership service, or gossip it between nodes, so every router sees the same ring. They also use a fast non-cryptographic hash, since nobody here is an attacker. The MD5 prefix is only there to keep the example to the standard library.

07

Field guide: who uses it and why


SystemWhat sits on the ringWhy consistent hashing
Amazon DynamoStorage nodes, by key rangeIn a large fleet, nodes join and leave all the time. Each change only moves data between ring neighbours, and virtual nodes let bigger machines take a bigger share.
Apache CassandraCluster nodes, by token rangePartitions data across the cluster. A new node streams one range from its neighbours instead of reshuffling the whole dataset.
DiscordChat servers, by session or guildLong-lived, stateful sessions. When a node changes, only the sessions in its arc have to reconnect somewhere else.
Akamai CDNCache servers, by content URLKeeps caches warm as servers come and go. This is the original use case from the Karger et al. work at MIT.
Google MaglevBackends behind a network load balancerWhen the set of backends changes, packets from an existing connection keep going to the same backend, so connections don't break.
Systems the book cites as using consistent hashing, and why it fits each one.

What they share: it costs something every time a key changes owner (a cold cache, a copied partition, a dropped connection), and servers come and go too often to stop everything and remap.

08

Interview ears: the follow-ups that come next


Interview tip"Why not just hash % N?" Changing N remaps nearly every key. In a cache, that means a storm of misses hitting the database during a failure or a scale-up.
Interview tip"A node dies. What happens?" Its range moves to the next node clockwise. Without virtual nodes, that one node takes the whole load. With them, the load spreads across many survivors.
Interview tip"How many virtual nodes?" Enough to bring variance under your target: 100 to 200 gives roughly 10% down to 5% standard deviation. You pay in routing metadata and membership churn.
Interview tip"Does it fix hot keys?" It reduces hot shards by spreading keys evenly. A single hot key still has one owner, so you replicate it or cache it closer to the client.
Interview tip"How do you find what to migrate on a join?" Walk anticlockwise from the new node to the node before it. Only that range moves, and it all comes from one existing node.
Q&A

Check yourself


Q1A 10-node cache that places keys with hash % N loses one node at peak. What do you expect to see first?
  • Only the dead node's tenth of the keys miss
  • Most keys miss at once and database load spikes
  • No misses, because hashes are deterministic
✓ Most keys miss at once and database load spikes — Dividing by 9 instead of 10 changes most remainders, so most keys now point at the wrong server.
Q2On a plain ring with one point per server, a server fails. Which server is most at risk of overload next?
  • The next server clockwise, which takes the whole range
  • Every server equally, since the load spreads evenly
  • The server before it, anticlockwise
✓ The next server clockwise, which takes the whole range — Keys walk clockwise to the next server, so one neighbour inherits the whole range. In the playground, zulu went from 4 keys to 8.
Q3You run mixed hardware and some nodes have twice the RAM. What is the cleanest way to give them more keys?
  • Hash the big nodes' names with a different function
  • Give the big nodes about twice as many virtual nodes
  • Use hash % N and list the big nodes twice
✓ Give the big nodes about twice as many virtual nodes — A server's share of the ring grows with its number of points, and the ring keeps its small-remap guarantee.
Q4Traffic to one viral key is melting its server. Will raising virtual nodes from 100 to 200 help?
  • Yes, the variance halves from 10% to 5%
  • Yes, the key gets split across two servers
  • No, one key still has exactly one owner
✓ No, one key still has exactly one owner — Virtual nodes balance how many keys each server holds, not the traffic to one key. Replicate or cache that key instead.
Sources: Book, Ch. 5: Design Consistent Hashing: rehashing problem, hash ring, server lookup, virtual nodes, affected keys, real-world uses (pp. 71-86)