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.
Outcomes
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.
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.
Hash every server onto the ring once, when it joins.
Hash the key with the same function.
Walk clockwise from the key's position. The first server you reach owns the key.
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.
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.
Press a button to make a change.
hash % NHow many keys plain modulo hashing would have relocated.
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.
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.
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.
Nothing outside that range moves. So rebalancing is a range copy between two neighbours, not a reshuffle of the whole cluster.
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.
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.
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.
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: TrueThe 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.
| System | What sits on the ring | Why consistent hashing |
|---|---|---|
| Amazon Dynamo | Storage nodes, by key range | In 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 Cassandra | Cluster nodes, by token range | Partitions data across the cluster. A new node streams one range from its neighbours instead of reshuffling the whole dataset. |
| Discord | Chat servers, by session or guild | Long-lived, stateful sessions. When a node changes, only the sessions in its arc have to reconnect somewhere else. |
| Akamai CDN | Cache servers, by content URL | Keeps caches warm as servers come and go. This is the original use case from the Karger et al. work at MIT. |
| Google Maglev | Backends behind a network load balancer | When the set of backends changes, packets from an existing connection keep going to the same backend, so connections don't break. |
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.