ulearn/systems

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:

  1. Place each server on the ring by hashing its name.
  2. Place each key on the same ring with the same hash.
  3. 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 movesRing movesIdeal (1 server's share)
Kill 1 of 367.3%35.6%33.3%
Add 1 to 375.0%26.8%25.0%
Kill 1 of 579.9%21.0%20.0%
Add 1 to 583.5%18.5%16.7%
Kill 1 of 1090.3%10.7%10.0%
Add 1 to 1091.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 each70.0%2.24× fair
3 v-nodes each42.3%1.62× fair
10 v-nodes each29.3%1.32× fair
50 v-nodes each5.9%1.07× fair
100 v-nodes each9.7%1.07× fair
150 v-nodes each8.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 moduloConsistent hashingRendezvous (HRW)Range partitioning
How a key finds its nodeservers[hash % N]Binary search for the next v-node clockwiseScore hash(key, server) for every server, take the highestLook 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 nodeOnly the ranges you choose to move
Lookup costO(1)O(log(N·v))O(N) hashes per keyO(log R) over R ranges, plus fetching the table
BalancePerfectNeeds ~100+ v-nodes per nodeGood with no tuningNeeds active splitting of hot ranges
Shared stateJust NServer list (the ring is derived)Just the server listA metadata service everyone reads
Range scansNoNoNoYes, a range lives on one node
Typical useFixed-size clusters, Kafka's default partitionerMemcached clients, Cassandra, Riak, DynamoCDN and proxy selection, small clustersBigtable, 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.