Consistent Hashing, Explained

After reading this you will know why hash mod N throws away almost every cache entry when you add one server, and why the ring moves only about 1/N of your keys instead.

What problem consistent hashing solves

You have a set of keys (user sessions, cache entries, database rows) and a set of servers to hold them. You need a rule that maps each key to exactly one server, and you want that rule to survive the day you add or remove a server.

The obvious rule is server = hash(key) mod N, where N is the server count. It spreads keys evenly and costs one modulo. It also has a fatal flaw: change N and nearly every key changes its answer.

Here is the hook. Suppose you run 4 cache servers and add a 5th. With hash mod N you switch from mod 4 to mod 5. A key with hash h lands on h mod 4 before and h mod 5 after. Those two values agree only when the change of divisor happens not to move that particular key. Count it over many keys and about 80% of them move. Every one of those keys is now looked up on a server that never cached it, so 80% of your traffic misses the cache at the same instant.

Consistent hashing changes the rule so that adding the 5th server moves only about 1/5 of the keys, and the other 4/5 stay exactly where they were.

The ring and the clockwise rule

Pick a hash function with a fixed output range. The common choice is a 32-bit space, so hashes run from 0 to 2^{32} - 1 (that is 4294967295). Bend that range into a circle so that the position after 2^{32} - 1 wraps back to 0.

Now place both servers and keys on the same circle:

  1. Hash each server's name (say "cache-a") to a position on the ring.
  2. Hash each key to a position on the same ring.
  3. A key belongs to the first server you meet when you walk clockwise from the key's position.

That is the whole rule. The server owns the arc of the ring from its predecessor server up to itself. When you add a server, it drops onto some point of the circle and takes over only the arc between itself and the previous server clockwise-behind it. Every key in that arc used to belong to the next server clockwise; now it belongs to the newcomer. No other key moves, because no other arc changed hands.

The 32-bit range is a convention, not a requirement. Cassandra's default partitioner uses a 128-bit Murmur3 token space. The math is identical; only the number of positions on the circle changes.

Why only K/N keys move

Assume keys are spread uniformly around the ring, which a good hash function gives you. With N servers placed on the circle, each server owns on average 1/N of the circumference. If you hold K keys total, each server owns about K/N keys.

\text{keys moved when adding one node} \approx \frac{K}{N+1}

Here K is the total key count, and N is the number of servers before the addition, so N+1 is the count afterward. The new node claims one arc of the enlarged ring, and that arc holds roughly K/(N+1) keys. Every other key keeps its owner.

Compare the two approaches directly. With hash mod N the fraction of keys that keep their server when you go from N to N+1 is only about 1/(N+1), so the fraction that moves is N/(N+1). With the ring the fraction that moves is 1/(N+1). Those are mirror images.

Fraction of keys that move when adding one server
Changehash mod N movesring moves
1 → 250%50%
2 → 366.7%33.3%
4 → 580%20%
9 → 1090%10%
99 → 10099%1%

At small N the two are close. The gap widens fast. At 100 servers, hash mod N reshuffles 99 keys out of every 100 for a single addition, while the ring touches 1.

The lumpiness problem and virtual nodes

The ring keeps its promise about movement, but a naive ring with one point per server is unfair. Random placement of a few points on a circle produces uneven gaps. With 4 servers each hashed to a single point, one server can easily own 40% of the circle while another owns 12%. Load then follows those arc widths, so your busiest server does three times the work of your quietest.

The fix is to give each server many points on the ring instead of one. These are virtual nodes (or vnodes). You hash "cache-a#0", "cache-a#1", up to "cache-a#49" and drop all 50 points onto the circle, all owned by cache-a. With more points per server, the total arc each server owns is the sum of many small gaps, and sums of many random pieces cluster tightly around the mean.

The standard deviation of each server's load share falls roughly like 1/\sqrt{v}, where v is the number of virtual nodes per server. Going from 1 point to 25 points cuts the spread by about a factor of 5. That is why memcached clients and Cassandra use dozens to hundreds of points per server rather than one.

Virtual nodes solve fairness, not movement. Adding a server still moves about K/(N+1) keys total. Those keys now arrive from many donor servers in small batches rather than from one neighbor, which spreads the rebuild cost.

Reproducing the demo: 4 servers, 1000 keys

Load the tool with its defaults and follow the numbers.

  1. Start with 4 servers and K = 1000 keys. Each server owns about 1000/4 = 250 keys on a well-balanced ring.
  2. Add a 5th server. The ring formula predicts 1000/5 = 200 keys move. The tool highlights roughly 200 keys, and the four old servers settle near 200 each.
  3. Switch on the hash mod N comparison for the same 4 → 5 change. About 1000 \times 4/5 = 800 keys change owner. The tool highlights close to 800.
  4. Set virtual nodes to 1 and read the keys-per-server table. Expect a wide spread, for example one server near 430 and another near 120, with a standard deviation around 130.
  5. Slide virtual nodes to 25. The same table tightens toward 200 each, and the standard deviation drops to roughly 40 to 60.

The exact highlighted counts jitter by a few percent because placement is random, but the ratios (20% moved on the ring, 80% moved by hash mod N) hold every time.

The ring's moved-key count falls as you add servers; hash mod N climbs toward moving everything.

With 1 virtual node per server across 4 servers, ownership shares might be 43%, 25%, 20% and 12%, a standard deviation near 13 percentage points. At 25 virtual nodes the shares cluster near 25% each with a standard deviation near 5 points. At 100 they sit near 25% with a spread near 2.5 points.

When to use it, and when not to

Consistent hashing pays off when both of these hold: your server set changes over time, and moving a key is expensive (a cache miss, a data copy, a re-replication). Distributed caches, sharded key-value stores and content-delivery routing all fit.

Skip it when the mapping is cheap to rebuild or the server set is fixed. If you have exactly 3 shards forever and rebalancing is a rare planned event, plain hash mod 3 is simpler and spreads load perfectly. Consistent hashing trades some balance and a lookup structure (a sorted array of ring positions plus a binary search) for stability under change.

Consistent hashing balances keys, not request rates. If one key gets 100 times the traffic of the rest (a celebrity account, a viral video), the ring will still park it on a single server and that server will melt. Solve hot keys separately with replication or client-side caching.

Common mistakes

Using too few virtual nodes
One point per server can hand a single server 40% of the ring. Use at least 20 to 50 points per server; hundreds if your servers differ in capacity.
Hashing the wrong server identifier
If you hash a server's IP address and the IP changes on restart, that server's points jump around the ring and drag keys with them. Hash a stable name or explicit token instead.
Assuming keys are uniform when they are not
The K/N promise assumes hashes spread evenly. Prefix collisions or a weak hash break that. Use a good hash (Murmur3, xxHash, or a cryptographic hash if you need it), not Object.hashCode-style sums.
Forgetting weighted capacity
A server with twice the RAM should own twice the ring. Give it twice as many virtual nodes rather than trying to size arcs by hand.

Related tools

To see the underlying lookup structure, the sorted ring of positions is a hash table plus a binary search over the Data Structure Visualizer. Consistent hashing decides which cache holds a key; what that cache evicts is a separate policy you can drive in the Cache Replacement Simulator. The membership question, which servers even exist right now, is what the Raft Consensus Visualizer settles. And once a hot key lands on one node, the Tail Latency & Autoscaling Simulator shows how that node's p99 latency blows up past 80% utilization.

Frequently asked questions

Does consistent hashing guarantee exactly K/N keys move?

No, it guarantees about K/(N+1) on average. A single new node claims one random arc, so the actual moved count varies. With virtual nodes the variance shrinks and the count sits closer to the mean.

Why 2^32 positions specifically?

It is just a convenient 32-bit range that many client libraries chose. The number of positions does not affect the movement math as long as it is large enough that collisions are negligible. Cassandra uses a 128-bit space; the behavior is the same.

What happens when a server dies?

Its arc merges into the next server clockwise, so that neighbor inherits its keys. Only about K/N keys move, the same fraction as an addition. With virtual nodes those keys spread across many neighbors instead of dumping onto one.

Is this the same as sharding?

Sharding is the general idea of splitting data across servers. Consistent hashing is one sharding scheme, chosen because it minimizes reshuffling when the server set changes. Range partitioning and directory-based sharding are alternatives with different tradeoffs.

How many virtual nodes should I use?

Enough that the standard deviation of load shares is acceptable. Since spread falls like 1/\sqrt{v}, going from 1 to 100 points per server cuts imbalance by about a factor of 10. Common defaults land between 128 and 256 per server.