The short answer
Quick answer: Consistent hashing is a way of assigning keys to servers so that adding or removing a server moves only a small fraction of the keys. Both servers and keys are hashed onto the same circular number line, the hash ring. Each key belongs to the first server found moving clockwise from the key's position. When a server joins or leaves, only the keys in the arc next to it change owner. With N servers, about 1/N of the keys move, instead of nearly all of them, as happens with the simple hash(key) % N approach.
The problem with hash mod N
Suppose you spread cache keys across 4 servers:
server = hash(key) % 4
It distributes keys evenly, and every client computes the same answer. Now add a fifth server and the formula becomes hash(key) % 5.
| Key hash | % 4 | % 5 | Moved? |
|---|---|---|---|
| 10 | 2 | 0 | Yes |
| 11 | 3 | 1 | Yes |
| 12 | 0 | 2 | Yes |
| 13 | 1 | 3 | Yes |
| 20 | 0 | 0 | No |
Going from 4 to 5 servers reassigns about 80% of all keys. For a cache, that means nearly every lookup misses at once and the database behind it is flooded. For a storage system, it means copying most of the data across the network. And servers change all the time: scaling up, scaling down, crashing.
This is the same weakness noted in sharding explained.
The hash ring
Consistent hashing, introduced by David Karger and colleagues at MIT in 1997 for distributed web caching, changes what is hashed.
- Take the output range of a hash function, say 0 to 2^32 - 1, and imagine it bent into a circle.
- Hash each server (by name or address) to a point on the circle.
- Hash each key to a point on the same circle.
- A key belongs to the first server clockwise from its position.
S1 (hash 10)
/ \
S4 (hash 80) S2 (hash 35)
\ /
S3 (hash 60)
key at 20 -> S2 key at 50 -> S3
key at 70 -> S4 key at 90 -> wraps round to S1
Adding a server
Add S5 at position 45, between S2 and S3. Keys between 35 and 45 used to belong to S3; now they belong to S5. Nothing else moves.
Removing a server
If S2 fails, its keys (those between 10 and 35) pass to the next server clockwise, S3. Again, nothing else changes.
In both cases the disruption is limited to one arc of the ring, roughly 1/N of the keys. The Wikipedia article gives the formal analysis.
The problem with the basic ring
With only a few points on the ring, placement is lumpy:
- Uneven load. Random positions leave some servers with long arcs and others with short ones. One server might own 40% of the keys and another 10%.
- Unfair failover. When a server dies, its entire load lands on a single neighbour, which may then be overloaded too.
Virtual nodes
The fix is to place each physical server on the ring many times, at different positions, by hashing server1#1, server1#2, and so on. Each position is a virtual node.
With 100 to 200 virtual nodes per server:
- Load evens out. Each server owns many small arcs, and the totals average out.
- Failures spread. A dead server's arcs are scattered round the ring, so its load is shared by many servers instead of one.
- Capacity can be weighted. Give a machine with twice the capacity twice as many virtual nodes.
The cost is a bigger lookup table. Finding the owner of a key is a binary search over the sorted list of positions, which is fast.
Replication on the ring
Storage systems keep each key on several servers. A simple rule: store the key on the next N distinct physical servers clockwise from its position. If one fails, the others still have the data.
This is the design described in the Amazon Dynamo paper, which brought consistent hashing into mainstream database design. See how database replication works for the quorum rules that go with it.
Where it is used
- Distributed caches. Memcached client libraries use it so that adding a cache server does not wipe out the whole cache. See caching strategies.
- Distributed databases. Cassandra, DynamoDB and Riak partition data on a ring.
- CDNs. Mapping content to cache servers was the original use; Akamai was founded on it.
- Load balancers. Routing each user or URL to the same backend, with minimal disruption when backends change. See what a load balancer does.
- Sharded real-time services. Assigning chat rooms or users to servers.
Variations
| Technique | Idea | Notes |
|---|---|---|
| Rendezvous (highest random weight) hashing | For each key, score every server with hash(key, server) and pick the highest | No ring to maintain; simple; work per lookup grows with the number of servers |
| Jump consistent hash | A tiny algorithm that maps a key to one of N numbered buckets | Very fast and even; buckets can only be added or removed at the end |
| Maglev hashing | A precomputed lookup table | Used in Google's network load balancer |
| Bounded-load consistent hashing | Cap each server's load and overflow to the next | Prevents a hot server |
A different approach to the same problem is to use a fixed number of slots. Redis Cluster hashes each key into one of 16,384 slots and assigns slots to nodes. Rebalancing means moving whole slots, tracked in an explicit table, rather than recomputing positions on a ring.
What it does not solve
- Hot keys. If one key is extremely popular, consistent hashing sends all of its traffic to one server. Spreading keys evenly does not spread requests evenly.
- Moving the data. It tells you which keys must move, but something still has to copy them while the system stays online.
- Agreement on membership. Every client must use the same list of servers. If they disagree, they route the same key to different places. Membership is usually shared through a coordination service or a gossip protocol.
- Range queries. Hashing scatters neighbouring keys, so scanning a range means asking every server.
A sketch in code
import bisect, hashlib
def h(value):
return int(hashlib.md5(value.encode()).hexdigest(), 16)
class Ring:
def __init__(self, servers, vnodes=100):
self.ring = sorted(
(h(f"{s}#{i}"), s) for s in servers for i in range(vnodes)
)
self.points = [p for p, _ in self.ring]
def lookup(self, key):
i = bisect.bisect(self.points, h(key)) % len(self.ring)
return self.ring[i][1]
The hash function here only needs to spread values evenly; it is not being used for security. The general idea of hashing is covered in how hash maps work.
Frequently asked questions
What problem does consistent hashing solve?
It lets you add or remove servers in a hashed, distributed system while moving only a small share of the keys, instead of almost all of them.
What is a virtual node?
One of many positions on the hash ring assigned to a single physical server. Virtual nodes even out load and spread the effect of a failure.
How many keys move when a server is added?
About 1/N of them, where N is the new number of servers.
Is consistent hashing the same as sharding?
It is one method of deciding which shard holds which key. Sharding is the broader practice of splitting data across servers.
Conclusion
Consistent hashing replaces "divide by the number of servers", which breaks whenever that number changes, with "walk round the ring to the next server", which barely notices. Add virtual nodes for balance and replicas for safety, and you have the partitioning scheme behind many distributed caches and databases.
Related articles
- Sharding Explained: How Databases Scale Beyond One Machine
- How Hash Maps Achieve O(1) Lookups
- What Is a Load Balancer and How Does It Decide Where Traffic Goes?
- Caching Strategies: Write-Through, Write-Back, and Cache-Aside
