· 6 min read
Consistent hashing: how it works and why virtual nodes matter
Consistent hashing lets you add or remove cache servers while moving only a small share of keys. Here is how the ring, virtual nodes and hash slots work.

Consistent hashing is a way to assign keys to servers so that adding or removing a server only moves a small share of the keys. With plain hash(key) % N, growing a cache cluster from 4 to 5 nodes remaps about 80% of keys; with consistent hashing and virtual nodes, only about 20% move, and all of them move to the new node. That difference decides whether scaling a cache is a non-event or a flood of misses against your database.
The idea is old and simple, and it sits underneath a lot of distributed caches, databases and load balancers. Once you see how the ring works, the trade-offs around virtual nodes and fixed slots make sense.
Why does hash modulo N break when you add a server?
The obvious way to spread keys across N cache servers is to hash the key and take the remainder:
server = hash(key) % N
It distributes evenly and costs nothing to compute. The problem shows up the moment N changes. A key whose hash is 17 lives on server 1 when N is 4, and on server 2 when N is 5. Almost every key gets a different remainder.
Running this over 100,000 keys confirms it: going from 4 servers to 5 moved 79.9% of them. In general, adding one server to N remaps roughly N / (N + 1) of all keys, so the bigger the cluster, the worse a single change gets.
For a cache, a moved key is a miss. Four out of five reads suddenly fall through to the database at the exact moment you added capacity, usually because load was already high. That is how a scaling action turns into the kind of cache stampede that takes down a database.
What is consistent hashing?
Consistent hashing places both servers and keys on the same circular number line, called the ring. Picture the full range of a hash function, say 0 to 2^64, bent into a circle so the largest value wraps around to zero.
Each server is hashed onto the ring by its name. Each key is hashed onto the ring too. To find a key's server, start at the key's position and walk clockwise until you hit a server. That server owns the key.
import bisect, hashlib
def h(s):
return int(hashlib.md5(s.encode()).hexdigest()[:16], 16)
class Ring:
def __init__(self, nodes):
self.points = sorted((h(n), n) for n in nodes)
self.hashes = [p[0] for p in self.points]
def get(self, key):
i = bisect.bisect(self.hashes, h(key)) % len(self.points)
return self.points[i][1]
Lookup is a binary search over the sorted server positions, so it stays cheap even with thousands of entries.
Now add a fifth server. It lands at one point on the ring and takes over only the arc between itself and the previous server going counter-clockwise. Every other key still walks clockwise to the same server it did before. Remove a server and only its keys move, to the next server along the ring. Nothing else is disturbed.
How do virtual nodes fix uneven load?
With one point per server, the arcs between servers are random lengths. Some servers own a huge slice of the ring and others a sliver. In a test with 5 servers and one point each, the busiest server held 1.94 times the average number of keys, and the quietest held less than half.
The fix is virtual nodes: each physical server is hashed onto the ring many times, under names like cache-3#0, cache-3#1 and so on. Each virtual node owns a small arc, and the many small arcs average out.
class Ring:
def __init__(self, nodes, vnodes=100):
self.points = sorted(
(h(f"{n}#{i}"), n) for n in nodes for i in range(vnodes)
)
self.hashes = [p[0] for p in self.points]
Measured on the same 100,000 keys, going from 4 to 5 servers:
| Virtual nodes per server | Keys moved | Busiest server vs average | Quietest vs average |
|---|---|---|---|
| 1 | 38.9% | 1.94x | 0.48x |
| 10 | 29.2% | 1.46x | 0.55x |
| 100 | 23.0% | 1.15x | 0.90x |
| 200 | 22.2% | 1.11x | 0.95x |
The ideal for adding a fifth server is 20%: the new node takes its fair share and nothing else moves. With 100 to 200 virtual nodes you get close to that, and in every run, every moved key went to the new server. None shuffled between the old ones.
Virtual nodes also give you weighting for free. A server with twice the memory gets twice the virtual nodes and owns about twice the keys. And when a server dies, its many small arcs are picked up by many different neighbors, so its load spreads across the cluster instead of landing on one unlucky server.
Fixed hash slots: the other common approach
Some systems skip the continuous ring and use a fixed number of hash slots instead. The key space is split into a set number of slots, keys map to slots with a stable hash, and a separate table maps slots to servers. Redis Cluster works this way with 16,384 slots.
The math is the same in spirit: the key-to-slot mapping never changes, so rebalancing means moving whole slots between servers and updating the table. The difference is control. With slots, an operator or a rebalancer decides exactly which slots move and when, and the mapping is an explicit, inspectable table. With a ring, placement falls out of hash positions.
| Approach | Keys moved when adding 1 of N servers | Balance | Who decides placement |
|---|---|---|---|
hash % N | About N/(N+1) | Even | The formula |
| Ring, 1 point per server | About 1/(N+1), uneven | Poor | Hash positions |
| Ring with virtual nodes | About 1/(N+1) | Good | Hash positions |
| Fixed hash slots | Only the slots you reassign | Good | A slot table |
When should you use consistent hashing?
Use it when the client or a routing layer has to pick a server for a key, and that set of servers changes while the system is running:
- Client-side sharding of caches. Memcached clients commonly use consistent hashing so a node failure or addition does not flush the whole cache.
- Sticky load balancing. Routing the same user or session to the same backend keeps local caches warm, and a backend restart only reshuffles its own users.
- Partitioned data stores. Databases that spread rows across nodes by key need a mapping that tolerates nodes joining and leaving.
You do not need it when the server count never changes, or when a central coordinator already owns placement. And it does not solve hot keys. If one key gets a tenth of all traffic, every scheme sends it to one server. That needs replication of the hot key or a local cache in front, not a better hash.
Two practical details matter in production. Every client must build the same ring, which means the same hash function, the same virtual node naming and the same server list, or different clients will disagree about where a key lives. And when a node is added, plan for the misses its new arc will cause; staggering additions or pre-warming keeps the database comfortable.
Key takeaways
hash(key) % Nremaps most keys whenNchanges; adding a fifth server moves about 80%.- Consistent hashing puts servers and keys on a ring and moves only about 1/(N+1) of keys per change.
- Virtual nodes even out load and spread a failed server's keys across many neighbors.
- Fixed hash slots give the same stability with explicit, operator-controlled placement.
- Neither approach fixes a single hot key; that needs replication or local caching.
FAQ
How many virtual nodes should each server have?
Somewhere around 100 to 200 per server gives good balance for small and medium clusters. More virtual nodes improve balance slightly but make the ring larger, which matters only for very large clusters.
Does consistent hashing guarantee perfectly even load?
No. It keeps movement small and, with virtual nodes, gets load close to even, but key popularity is not uniform. Uneven traffic per key still needs separate handling.
Which hash function should I use?
Any fast, well-distributed, non-cryptographic hash works, such as xxHash or MurmurHash. What matters most is that every client uses the same one, with the same server names and virtual node scheme.