Fundamentals · 12
Consistent hashing
Map keys to nodes so that adding or removing a node moves only a small fraction of keys, using a hash ring and virtual nodes.
Suppose you shard a cache across 4 servers with server = hash(key) mod 4. Add a fifth server and the formula becomes mod 5, and about 80% of keys now map to a different server. For a cache that means a sudden wave of misses that can take down the database. Consistent hashing fixes this.
The hash ring
Imagine the hash output space (say 0 to 2³²−1) bent into a circle.
- Hash each node (for example its IP) to a point on the ring.
- Hash each key to a point on the ring.
- A key belongs to the first node clockwise from it.
When node D joins, it takes over only the keys between its predecessor and itself, which used to belong to the next node clockwise. Every other key stays where it was. When a node leaves, its keys move to the next node clockwise. On average only K/N keys move (K keys, N nodes), instead of almost all of them.
Virtual nodes
With only a few points on the ring, the arcs between nodes are uneven, so some nodes get far more keys than others. And when a node dies, all its load lands on one neighbour.
Virtual nodes fix both: each physical server is hashed to many points (say 100–200), named A#1, A#2 and so on.
- Load evens out statistically across servers.
- A dead server’s keys spread across many other servers, not one.
- Stronger machines can take more virtual nodes, which works as weighting.
Where it’s used
- Distributed caches: Memcached client libraries, Redis Cluster (with 16,384 fixed hash slots, a close relative).
- Partitioned databases: Cassandra, DynamoDB, Riak place data and replicas on a ring. Replicas go on the next N distinct nodes clockwise.
- Load balancers: route requests for the same key to the same backend, which keeps caches warm.
- CDNs and object stores: decide which server holds which object.
Pros
- Adding or removing nodes moves only ~1/N of keys
- No central lookup table needed; any client can compute the owner
- Virtual nodes give even load and smooth failure handling
Cons
- Needs virtual nodes to balance well
- Range queries aren’t possible (keys are scattered by hash)
- Hot keys still overload their single owner
- Membership changes must reach every client (gossip or config service)
Test yourself
Answer in your head, then click a card to check. All cards are in the Anki deck.