Consistent Hashing
Add a server and move only 1/N of the keys: why hash mod N reshuffles everything, how a hash ring limits a change to one arc, how virtual nodes even out the load, and what no hashing scheme can do about a hot key.
An interactive System Design lesson: 19 steps, about 30 minutes, on a live simulation in your browser.
One cache server cannot hold every user profile, so the app spreads them over four: cache-a to cache-d. Every key must live in exactly one place, and anyone who wants it must be able to work out where without asking around.
The obvious rule: hash the key to a big number and take it modulo the number of servers. hash(user:7) mod 4 = 0, so user:7 lives on cache-a, the first in the list. The router computes this for every request; so could every app server, with no shared table at all.
What you will learn
The obvious rule: hash mod N
- Four caches, one rule: A sharding rule is a function from key to server. Anyone who knows the function and the server list can find any key without a lookup table.
- Add one cache: Under hash mod N, changing N changes the answer for almost every key. The fraction that stays is about 1/N, not the fraction that moves.
- Take it out again: Every key that changes owner is a cache miss or a data migration. A rule that moves 90% of keys on any membership change makes scaling and failure equally expensive.
The ring
- Keys and servers on one circle: On a hash ring each server owns the arc that ends at its position. A key's owner depends only on the nearest server clockwise, not on how many servers there are.
- Drill: read the ring
- Add a server to the ring: Adding a server to a ring moves only the keys on the arc it lands in, and they all move to the new server. Nothing shuffles between the old servers.
- A server leaves the ring: When a server leaves a ring its whole arc goes to the next server clockwise. Few keys move, but they all land on one neighbour.
Virtual nodes
- One server owns over half the keys: With one point per server, each server's share is the length of a random arc. The ring limits movement; it does nothing for balance.
- Sixteen points per server: Virtual nodes turn one random arc per server into many small ones. The more points per server, the closer every server gets to 1/N of the keys.
- Adding a server takes a fair share: With virtual nodes, a new server takes about 1/N of the keys, a little from every existing server, and nothing else moves.
- Removing one spreads its keys
- What it looks like in real config
- Drill: how much moves?
Hot keys
- One key everybody wants: Hashing balances keys, not traffic. A hot key lands on one server whatever the ring looks like.
- Cooling a hot key: A hot key is fixed by making copies of it (in the app, under several names, or on several nodes), because a hash function can only ever name one owner.
Where you meet it
- Where you meet a ring: Anywhere a changing set of servers must agree on who owns a key without a central table, you will find a ring or its cousin, fixed hash slots.
- Replicas: the next servers clockwise: On a ring, replication is placement: a key's N copies live on the N distinct servers clockwise from it.
Recap & playground
- Cheat sheet
- Playground