A cache tier of N nodes sits behind a proxy that picks a node with hash(key) mod N. What happens to the hit rate when one node is removed, and how does consistent hashing change the outcome?
answer
- the divisor is the fleet size
- change N, change every answer
- nodes and keys on one circle
- only the departing arc moves
- many positions per physical node
basics
~20 sModulo hashing remaps most keys when N changes, so nearly the whole cache misses at once and the origin takes the load. Consistent hashing maps keys and nodes onto a ring, so removing a node moves only that node's share — roughly 1/N of keys.
solid answer
~60 sWith `hash(key) mod N`, the node index depends on N itself, so changing N from 4 to 3 changes the answer for about three quarters of all keys. Every one of those keys is now looked up on a node that never stored it: a near-total cache miss storm, arriving at your origin all at once, at the exact moment you have one node fewer. Consistent hashing removes N from the arithmetic. Nodes and keys are hashed onto the same circular keyspace, and a key belongs to the first node clockwise from it. Removing a node hands only its arc to its successor — about 1/N of keys move, and every other key keeps the node it had. Each node is placed at many points on the ring (virtual nodes), otherwise arcs come out wildly uneven and the departing node dumps its entire share on one successor. The trade you accept is that routing is now by key identity rather than by load, so a hot key still lands on one node.
code
python · 23 linesimport bisect, hashlib
def h(s):
return int(hashlib.md5(s.encode()).hexdigest()[:8], 16)
keys = [f"key{i}" for i in range(10000)]
nodes = [f"node{i}" for i in range(4)]
# modulo: the divisor changes, so most keys change owner
moved = sum(1 for k in keys if h(k) % 4 != h(k) % 3)
print("modulo remapped:", moved / len(keys))
def ring(ns, vnodes=100):
r = sorted((h(f"{n}#{v}"), n) for n in ns for v in range(vnodes))
return [p[0] for p in r], [p[1] for p in r]
def owner(pos, names, k):
return names[bisect.bisect(pos, h(k)) % len(pos)]
p4, n4 = ring(nodes)
p3, n3 = ring(nodes[:3])
moved = sum(1 for k in keys if owner(p4, n4, k) != owner(p3, n3, k))
print("consistent-hash remapped:", moved / len(keys))go deeper
Know that hashing a key to a node gives you locality — the same key always reaches the same cache — and that plain modulo hashing breaks that as soon as the number of nodes changes.
Explain the ring: nodes and keys share a keyspace, a key belongs to the next node clockwise, and only the departing node's arc changes hands, which is roughly 1/N of keys.
Reason about the incident shape: a resize under modulo hashing throws a full miss storm at an origin that is already degraded, and virtual nodes are what stop a departure from toppling a single successor.
Own the tension between locality and load-awareness across the platform: when key-affinity routing is worth its hot-key risk, whether membership needs consensus rather than gossip, and what the blast radius of a ring change is.
## Why modulo hashing is fragile `hash(key) mod N` is attractive because it is stateless: any proxy can compute the owner with no coordination. The flaw is that the divisor is the fleet size. Change N and you change the arithmetic for every key at once. Going from four nodes to three, only keys whose hash happens to give the same index under both moduli stay put — empirically about a quarter of them; roughly three quarters move. Every moved key is a guaranteed miss on the new owner, plus a stranded copy on the old one. The operational shape of that is worse than the number suggests. Node loss is usually not a calm event: you lose a node because it crashed or was drained, and the cache tier exists precisely because the origin cannot serve full traffic. Now the origin receives a near-total miss storm while it is already unhappy, every node begins filling with keys it does not have, and memory pressure evicts keys that were still valid. Adding a node is the same event with the same cost — which is why teams with modulo hashing quietly stop resizing their cache tier at all. ## The ring Consistent hashing puts nodes and keys into one circular keyspace, for example the range of a 32-bit hash. Each node is hashed to a position on the circle; each key is hashed to a position too, and belongs to the first node encountered walking clockwise. Lookup is a binary search over the sorted node positions. The crucial property: a node's ownership is defined by *its neighbours on the ring*, not by the total count. Remove a node and only the arc between it and its predecessor changes hands, passing to the next node clockwise. Every other key's answer is unchanged. Add a node and it takes a slice out of exactly one successor. The share that moves is about 1/N — the minimum possible, since those keys had to move somewhere. ```python import bisect, hashlib def h(s): return int(hashlib.md5(s.encode()).hexdigest()[:8], 16) # one node at one ring position gives very uneven arcs; # hashing each node under many labels evens them out positions = sorted(h(f"node{i}#{v}") for i in range(4) for v in range(100)) print(bisect.bisect(positions, h("cart:42")) % len(positions)) ``` ## Virtual nodes, and why a naive ring is not enough With one ring position per node, arc sizes are the gaps between random points: highly variable, so one node can own several times another's share. Worse, when a node leaves, its *entire* arc lands on a single successor — which may already be the largest owner — so removing one node can double another's load and take it down too. The fix is virtual nodes: hash each physical node under many labels (`node2#0`, `node2#1`, …), typically 100–200 positions per node. Arcs become numerous and small, so the law of large numbers evens out each physical node's total share, and a departing node's many small arcs are spread across all remaining nodes instead of one. Virtual nodes also give you weighting for free — a machine with twice the memory simply gets twice as many positions. ## What consistent hashing buys, and what it costs It buys **locality**: the same key reliably reaches the same node, so per-node caches stay warm, an in-memory index or shard stays on the node that owns it, and rate-limit or session counters can be kept locally without a shared store. It buys **graceful membership change**: scale-out, scale-in and node loss cost about 1/N of the cache instead of all of it. What it costs is load-awareness. Routing is now a function of the key, not of how busy the target is, so a genuinely hot key — one product on the front page, one large tenant — saturates its owner while the rest of the ring idles, and no amount of balancing helps because every request for that key is required to go to that node. Mitigations exist (replicate hot keys across a few successors, add a bounded-load rule that spills to the next node when the owner is over a threshold, or fall back to a load-based algorithm for keys detected as hot), and interviewers like to hear that you know hashing and load-awareness are in tension rather than a free upgrade. The other cost is state: every proxy must agree on the ring, so membership changes have to propagate. Two proxies with different views route the same key differently, which for a cache means a duplicated entry and for a stateful shard can mean a correctness problem — which is why sharded stateful systems put ring membership behind a consensus mechanism rather than letting each proxy guess.
- Why do implementations place each node at a hundred or more ring positions instead of one?With one position per node, arcs are random gaps and come out very uneven, and a departing node dumps its whole arc on a single successor — potentially toppling it. Many small positions per node even out each node's total share and spread a departure across all survivors. They also give weighting: a bigger machine simply gets more positions.
- Consistent hashing is in place and one node is still saturated. What is the likely cause?A hot key. Hashing routes by key identity, so all traffic for the single hottest key is required to reach one node no matter how idle the others are. Fixes are replicating hot keys to the next few nodes on the ring, a bounded-load rule that spills over a threshold, or detecting hot keys and serving them load-balanced.
- Two proxies in front of the same ring disagree about membership. What goes wrong?They compute different owners for the same key. For a cache that means duplicate entries and a lower hit rate — wasteful but safe. For a stateful shard it can mean two nodes accepting writes for the same key, which is a correctness problem; that is why sharded stores put ring membership behind a consensus mechanism rather than letting each client guess.
saying these in an interview costs you the question
- Thinks mod N only remaps the keys of the lost node
- Believes consistent hashing prevents cache misses entirely
- Omits virtual nodes and assumes arcs are even
- Expects hashing to balance load as well as it balances keys
- Assumes all proxies see membership changes instantly