topics / data-caching
Consistent Hashing
Spread a cache over N servers with hash(key) % N and it works, until N changes. Kill one node and watch nearly every key move, then do the same on a hash ring and watch only the dead node's keys go.
The problem: the modulo cache stampede
You run a cache in front of a database, and one machine isn't enough, so you spread keys over N cache servers. The obvious way to pick a server is servers[hash(key) % N]: fast, stateless, and perfectly even. Every client computes the same answer with no coordination.
Then a server dies, and N goes from 5 to 4. The remainder of hash % 4 has almost nothing to do with hash % 5, so it isn't just the dead server's keys that move. Across 10,000 keys, killing one of five servers sends 79.9% of them to a different server. The dead server only held about 20%; the other ~60% were sitting on healthy machines, and now every client looks for them somewhere else.
For a cache, “somewhere else” means a miss. Most of the cache goes cold in one step, and every miss falls through to the database at the same moment. The database was sized for the trickle of misses a warm cache lets through, not for that: this is a cache stampede, and it turns one lost cache node into a database outage. Adding a server to relieve load does the same thing, which is worse: scaling up under pressure knocks over the database.
The Hashing & Collisions topic shows the in-memory version of this: resizing a hash table rehashes every key. That's fine when the buckets are array slots. It isn't fine when they're machines.
The hash ring
Consistent hashing (Karger et al., 1997) stops the server count from appearing in the formula at all. Take the whole hash space, 0 to 2³²−1, and bend it into a circle so the top wraps back to 0. Then:
- Place each server on the ring by hashing its name.
- Place each key on the same ring with the same hash.
- A key belongs to the first server clockwise from it.
// ring: every virtual node's position, sorted ascending
function ownerOf(key: string): Node {
const h = hash(key); // 0 … 2³² − 1
let i = lowerBound(ring, h); // first v-node with pos >= h
if (i === ring.length) i = 0; // walked off the end: wrap to the start
return ring[i].node;
}Each server owns the arc between its predecessor and itself. When a server dies, its arc merges into the next server clockwise: its keys slide to that one neighbour, and every other key's “first server clockwise” is unchanged. When a server joins, it lands inside one existing arc and takes only the part before it. Nothing else anywhere on the ring moves.
The math of minimal disruption
The floor: K / N
With K keys on N servers, a server holds about K / N of them. When it dies, those keys must move: their home is gone. So K/N is the minimum any scheme can achieve, and it's exactly what a ring moves. Adding a server is symmetric: it needs K / (N + 1) keys to carry its share, and that's all it takes.
Modulo: K · (N − 1) / N
A key stays put under % N → % (N ± 1) only when both remainders pick the same server, which for a well-spread hash happens about 1 time in N. So roughly (N − 1) / N of all keys move, and it gets worse as the cluster grows: 75% at 4 nodes, 90% at 10, 99% at 100.
Measured over 10,000 keys with the simulation's own hash, and 100 virtual nodes per server on the ring:
| hash % N moves | Ring moves | Ideal (1 server's share) | |
|---|---|---|---|
| Kill 1 of 3 | 67.3% | 35.6% | 33.3% |
| Add 1 to 3 | 75.0% | 26.8% | 25.0% |
| Kill 1 of 5 | 79.9% | 21.0% | 20.0% |
| Add 1 to 5 | 83.5% | 18.5% | 16.7% |
| Kill 1 of 10 | 90.3% | 10.7% | 10.0% |
| Add 1 to 10 | 91.1% | 8.9% | 9.1% |
The ring doesn't hit the ideal exactly, because a real server's share is only about 1/N. How close it gets depends on the next section.
Virtual nodes and hot spots
With one position per server, arc lengths are luck. Five random points on a circle almost never split it into fifths: one server usually lands after a long empty stretch and owns far more than its share, and when a server dies, its entire load lands on one neighbour, which may already be the busiest node.
The fix is virtual nodes: hash each server onto the ring many times ("N2#vn0", "N2#vn1", …). Each server now owns many small arcs scattered around the circle, and the sum of many random arcs is far more predictable than one. A dead server's keys also scatter across many successors instead of dumping onto one. For 5 servers:
| Load std-dev (vs fair share) | Busiest node's load | |
|---|---|---|
| 1 v-node each | 70.0% | 2.24× fair |
| 3 v-nodes each | 42.3% | 1.62× fair |
| 10 v-nodes each | 29.3% | 1.32× fair |
| 50 v-nodes each | 5.9% | 1.07× fair |
| 100 v-nodes each | 9.7% | 1.07× fair |
| 150 v-nodes each | 8.3% | 1.12× fair |
Going from 1 to 100 v-nodes takes the spread from 70.0% to 9.7%, and the busiest node from 2.24× to 1.07× its fair share. Spread shrinks roughly with 1/√v, so the curve flattens: most of the win is in the first 50–100, which is why memcached's ketama clients use 160 points per server and Cassandra historically used 256 tokens per node (16 by default since 4.0, paired with a smarter token allocator).
Virtual nodes also make weighting trivial: a server with twice the memory gets twice the v-nodes. The cost is memory and lookup time for the ring itself (a sorted array of N × v positions, binary searched), which is small.
Modulo vs ring vs rendezvous vs ranges
| Naive modulo | Consistent hashing | Rendezvous (HRW) | Range partitioning | |
|---|---|---|---|---|
| How a key finds its node | servers[hash % N] | Binary search for the next v-node clockwise | Score hash(key, server) for every server, take the highest | Look the key up in a range → node table |
| Keys moved when one node joins/leaves | ~(N−1)/N of all keys | ~1/N (one node's share) | ~1/N, and only to/from that node | Only the ranges you choose to move |
| Lookup cost | O(1) | O(log(N·v)) | O(N) hashes per key | O(log R) over R ranges, plus fetching the table |
| Balance | Perfect | Needs ~100+ v-nodes per node | Good with no tuning | Needs active splitting of hot ranges |
| Shared state | Just N | Server list (the ring is derived) | Just the server list | A metadata service everyone reads |
| Range scans | No | No | No | Yes, a range lives on one node |
| Typical use | Fixed-size clusters, Kafka's default partitioner | Memcached clients, Cassandra, Riak, Dynamo | CDN and proxy selection, small clusters | Bigtable, HBase, Spanner, CockroachDB, managed DynamoDB |
Rendezvous hashing (Thaler & Ravishankar, 1996) reaches the same minimal disruption without a ring: every client scores every server for a key and picks the winner. A dead server only loses the keys it was winning, and each of those goes to whichever server scored second for it, so they spread evenly with no virtual nodes. The catch is O(N) work per lookup, which is nothing for 10 servers and a lot for 10,000.
Range partitioning gives up on computing the owner and writes it down instead. That costs a metadata service, but it's the only option here that keeps sorted scans on one machine, and it lets the system split one hot range without touching any other.
In production: Discord and DynamoDB
Discord: routing guilds and sessions
Discord's real-time gateway is written in Elixir. Each guild (a Discord server) runs as a process on one node of a cluster, and every connected user has a session process holding their websocket. When a message is posted, the guild process fans it out to the sessions of everyone in it; when a session starts, it has to find the process for each of its guilds.
They find it with a consistent hash ring over the guild id, so any node can compute a guild's home without asking a central directory, and adding or removing a node only relocates the guilds on its arcs rather than every guild in the cluster. Discord open-sourced the ring as ex_hash_ring. At their scale the ring lookup itself sits on the hottest path in the system, which is why they tuned it for speed rather than treating it as a detail.
Amazon Dynamo → DynamoDB: storage shards
Amazon's 2007 Dynamo paper is the textbook consistent-hashing store. Every node takes many “tokens” (virtual nodes) on the ring; each key is stored on the first N distinct physical nodes clockwise from it, its preference list, so when a node dies its keys already have warm replicas on exactly the neighbours that inherit its arcs. The paper also reports that purely random tokens made rebalancing and bootstrapping slow, and they moved to splitting the ring into equal-sized partitions that are assigned to nodes as whole units.
The managed DynamoDB service took that last step further. It still hashes each item's partition key, but it assigns contiguous ranges of the hash space to partitions and records them in a routing metadata service. Partitions split when they grow too large or too hot, and move between hosts one at a time. That's range partitioning over hashed keys: the even spread of hashing plus the explicit control of a metadata table, paid for with a service that has to stay available.
Beyond the basic ring
- Replication fixes the cold handover. Even a perfect ring moves K/N keys on a failure, and for a plain cache those are still misses. Storing each key on the next R nodes clockwise means the node that inherits a dead node's arc already has the data.
- Bounded loads. Virtual nodes even out the hash space, not the traffic: one viral key is still one key. Consistent hashing with bounded loads (Mirrokni et al., Google, 2016) caps every server at, say, 1.25× the average and sends overflow to the next server clockwise. Vimeo added it to HAProxy for exactly this reason.
- Jump consistent hash. Lamping and Veach's 2014 algorithm needs no ring and no memory, just a few lines of arithmetic, and balances perfectly. It only works when servers are numbered 0…N−1 and can only be added or removed at the end, so it suits storage shards better than a cache fleet where any node can die.
- Changing the v-node count moves keys too. Raising it adds positions that each take a slice of a neighbour's arc, so going from 10 to 20 per node moves about half the keys, even though no server came or went. Pick the count once.