System designCore6 min3 things to try
Consistent hashing
How to add a server to a cluster so that only a sliver of the keys move, instead of almost all of them.
1 · The problem
Add one server, and almost everything moves
The obvious way to pick a server for a key is hash(key) % N. It spreads keys evenly, until N changes. Add a server and see where the keys land now.
Every moved key is a cache miss that falls through to the database, all at once. In a cache cluster that is a thundering herd, at exactly the moment you were adding capacity.
2 · The idea
Put keys and servers on the same ring
Hash the servers and the keys onto the same circle. A key belongs to the first server clockwise from it. When a server joins, it only takes the keys between itself and the server before it.
- S1 owns 7 keys
- S2 owns 13 keys
- S3 owns 4 keys
Ring
0
keys moved
Modulo, same change
0
keys moved
3 · How a lookup works
Finding a key's server
- Hash the key
to a position on the ring, a number between 0 and 2³².
- Walk clockwise
to the next server position. Clients keep the sorted list of server positions, so this is a binary search: O(log N).
- Wrap around
past the top: the first server on the ring owns the keys after the last one.
- Replicate
by also writing to the next one or two distinct servers clockwise. That is how a database keeps a copy when a node dies.
4 · The catch
Few servers make lopsided arcs: use virtual nodes
With one position per server, the gaps between them are random, so one server can own twice its share. Give each server many positions on the ring and the load evens out.
5 · In a real system
Where the URL Shortener uses it
URL Shortener System Design
Cassandra places every short link on a token ring
Each Cassandra node owns ranges of the ring, with several virtual nodes each, and a link is stored on its owner plus the next replicas clockwise. Adding a node takes over a share of the ring from its neighbours; the rest of the links don't move, so the redirect path keeps working while the cluster grows.
6 · When to use it
And what else there is
| Approach | Good at | Watch out for |
|---|---|---|
| Consistent hashing | Caches and stores where nodes come and go | Needs virtual nodes for even load; range queries are scattered |
| Modulo hashing | A fixed number of servers | Almost every key moves when the count changes |
| Rendezvous hashing | Even spread without virtual nodes | Each lookup scores every server: O(N) |
| Range partitioning | Range scans, such as everything created today | Hot ranges, and a rebalancer to run |
Check yourself
3 questions
Takeaways
Remember this
- Modulo hashing moves almost every key when the server count changes; a ring moves about 1/N of them.
- A key belongs to the first server clockwise; replicas go to the next distinct servers.
- Virtual nodes even out the load and let a bigger machine take more positions.