Lesson 8 / 25
Consistent Hashing
Add servers without reshuffling everything.
A ring with virtual nodes
Assigning keys with hash(key) % N moves almost every key when N changes, causing a storm of cache misses or data moves. Consistent hashing places servers and keys on a ring; each key belongs to the next server clockwise, so adding a server moves only about 1/N of keys. Virtual nodes (many points per server) smooth out the load. It is used in distributed caches, Dynamo-style databases and partitioned queues.
Modulo versus a ring when adding a fifth server, run
I ran this with Python 3.12.3 using only the standard library; inputs are fixed or seeded, so the output is reproducible. With modulo hashing about 80% of 10,000 keys move; with a 100-virtual-node ring about 20% move, which is the minimum possible, and load stays roughly balanced.
# Modulo hashing vs consistent hashing when adding a 5th server
import bisect, hashlib
def h(s):
return int(hashlib.md5(s.encode()).hexdigest(), 16)
keys = [f"user:{i}" for i in range(10_000)]
def modulo(n):
return {k: h(k) % n for k in keys}
class Ring:
def __init__(self, nodes, vnodes=100):
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 node(self, key):
i = bisect.bisect(self.hashes, h(key)) % len(self.points)
return self.points[i][1]
before, after = modulo(4), modulo(5)
moved = sum(before[k] != after[k] for k in keys)
print(f"modulo 4 -> 5 servers : {moved / len(keys):.1%} of keys move")
r4, r5 = Ring(["s0", "s1", "s2", "s3"]), Ring(["s0", "s1", "s2", "s3", "s4"])
moved = sum(r4.node(k) != r5.node(k) for k in keys)
print(f"ring 4 -> 5 servers : {moved / len(keys):.1%} of keys move (ideal 20%)")
load = {}
for k in keys:
load[r5.node(k)] = load.get(r5.node(k), 0) + 1
print("keys per server:", dict(sorted(load.items())))
Output:
modulo 4 -> 5 servers : 79.6% of keys move
ring 4 -> 5 servers : 19.9% of keys move (ideal 20%)
keys per server: {'s0': 1818, 's1': 2069, 's2': 2006, 's3': 2122, 's4': 1985}Mention virtual nodes
Without them, a few servers can own much larger arcs of the ring and become hot.
Quick check: Roughly what fraction of keys moves when going from 4 to 5 servers with consistent hashing?
- All of them
- About 80%
- About 20%
- None
Answer
About 20% — Only the new server's share moves.