Consistent hashing
Mapping keys to nodes so that adding or removing a node moves only a small share of keys, instead of reshuffling almost all of them.
Distribution
Learn it
The obvious way to spread keys over N servers is
hash(key) % N. It balances well, until N changes. Going from 10 to 11 cache servers changes the server for most keys, so the hit rate collapses and every key is fetched from the database at once.Anything that routes by key (caches, connection servers, shards) needs a mapping that survives membership changes.
Consistent hashing places nodes and keys on the same circular hash space, a ring. A key belongs to the first node clockwise from its hash. Adding a node takes over only the keys between it and its predecessor; removing one hands its keys to the next. On average only 1/N of keys move.
Add a node under each placement, and compare one point per node with a hundred.
2,000 keys spread over cache nodes. Add a node and count the keys that now belong somewhere else: each one is a cache miss, or data to copy. 4 nodesKeys per node. The line is a perfectly even share.
- Node 1500
- Node 2508
- Node 3512
- Node 4480
The busiest node holds 1.0× an even share.
Check
What do virtual nodes (many ring points per physical node) improve?Consistent hashing decides where a key lives, not how much load it brings: a single hot key still lands on one node. Alternatives with the same goal include rendezvous hashing (each key picks the node with the highest hash(key, node)) and directory-based placement (a lookup table you can edit deliberately).
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Hot key
- One key's load exceeds a node; hashing cannot split it.
- Disagreeing membership
- Clients with different views of the ring send the same key to different nodes.
- Too few virtual nodes
- Uneven arcs give some nodes several times the load of others.
Instead, consider
- Modulo hashing
- The number of nodes never changes, or reshuffling everything is cheap.
- Directory / range map
- You need to move specific ranges deliberately (e.g. to rebalance hot shards) rather than by hash.
- Rendezvous hashing
- Node counts are small and you want simple, even placement without virtual nodes.
In practice
- Cassandra / ScyllaDB / DynamoDB
- Token rings with virtual nodes and replication.
- Client-side memcache routing
- Consistent hashing in the client library or a router such as mcrouter.
- Load balancer hash policies
- Route by a header or path so one key's requests reach one backend.
It assumes
- Clients (or a router) agree on the current membership of the ring.
- Keys are numerous and individually small compared with a node's capacity.
Explain it in your own words
Where you practise it
Further reading
Engineers describing it in systems they run.
- How Discord Stores Trillions of Messages
Discord · Bo Ingram · Post, Mar 2023
The same data model six years on: hot partitions, a service layer that merges identical reads, and a migration of the full history to a new database.
Related concepts
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
- Caching
Keeping a copy of data closer to where it is used, trading freshness and complexity for speed and reduced load on the source.
- Replication
Keeping copies of data on several machines for durability, read capacity and locality, and living with copies that briefly disagree.
- Request coalescing
When many callers ask for the same thing at the same moment, do the expensive work once and give every caller the result. The fix for thundering herds and hot keys.