AlgoMaster Logo

Consistent Hashing

High Priority11 min readUpdated September 16, 2026
AI Mock Interview

Practice this topic in a realistic system design interview

Listen to this chapter
Unlock Audio

Premium Video

This video is available to premium subscribers only

Unlock Full Access

Distributed systems often need to answer a simple but important question: which server should handle this key?

If that answer changes too often when servers join or leave, the system wastes time moving data.

Consistent hashing solves this by minimizing how many keys need to move when the cluster changes. This chapter covers how it works, the limitations of simple modulo hashing that it solves, and virtual nodes, which make it practical in production.

1. The Problem

Suppose our application caches user profiles across five servers, S0 through S4. When a request comes in, we need a reliable way to decide which server should handle that profile.

Sending requests to random servers would be inefficient, because the profile might already be cached on another server. Instead, we use the profile key to consistently map it to one server. As long as the server list stays the same, every request for the same profile goes to the same place.

The challenge is what happens when the cluster changes. If we add a server for more capacity or remove one after a failure, we want to update the routing without moving too many profiles around.

2. Modulo Hashing

Let's start with the simplest approach. We hash the key to get a number, then take that number modulo the number of servers. The result gives us the server index.

With five servers, the output is always between 0 and 4.

If a key hashes to 42, then 42 modulo 5 is 2, so the request goes to S2. If another key hashes to 17, 17 modulo 5 is also 2, so it goes to S2 as well.

Scroll
Key HashCalculationServer
4242 mod 5 = 2S2
1717 mod 5 = 2S2
1818 mod 5 = 3S3
3737 mod 5 = 2S2
9191 mod 5 = 1S1

With a good hash function, keys are distributed fairly evenly across the servers. Each server ends up owning roughly one fifth of the keys, and the same key always lands on the same server.

3. What Happens When Servers Change

Now suppose we add S5, bringing the total to six servers. The formula becomes hash(key) % 6, and the same keys are suddenly mapped differently.

42 modulo 6 is 0, so that key moves from S2 to S0. 17 modulo 6 is 5, so the other key moves from S2 to S5.

Scroll
Key HashWith 5 ServersWith 6 ServersMoved?
4242 mod 5 = 2, S242 mod 6 = 0, S0Yes
1717 mod 5 = 2, S217 mod 6 = 5, S5Yes
1818 mod 5 = 3, S318 mod 6 = 0, S0Yes
3737 mod 5 = 2, S237 mod 6 = 1, S1Yes
9191 mod 5 = 1, S191 mod 6 = 1, S1No

And this happens to many other keys as well. Even though we added just one server, a large portion of the cache gets reassigned. In the table, four of the five keys changed owners.

The cached data is still sitting on the old servers, so requests that hit the new owners miss the cache and have to fetch the data again.

Removing a server causes the same kind of reshuffling. If S4 fails and we drop to four servers, the formula becomes hash(key) % 4, and again most keys change owners.

Scroll
Key HashWith 5 ServersWith 4 ServersMoved?
42S242 mod 4 = 2, S2No
17S217 mod 4 = 1, S1Yes
18S318 mod 4 = 2, S2Yes
37S237 mod 4 = 1, S1Yes
91S191 mod 4 = 3, S3Yes

What we really want is to move only the keys that need to move, while keeping most existing cache entries useful.

4. The Hash Ring

Consistent hashing solves this by placing both servers and keys in a fixed hash space.

Loading simulation...

For simplicity, imagine positions from 0 to 999. After 999, we wrap back to 0, forming a ring. In real systems, the range is much larger, such as 0 to 2^32 - 1 or 2^64 - 1, but the idea stays the same.

We hash each server's identifier to place it somewhere on the ring. Suppose S0 lands at 100, S1 at 300, S2 at 500, S3 at 700, and S4 at 900.

These evenly spaced positions make the routing easier to follow. Each server owns the range of positions between the previous server and itself.

Scroll
ServerPositionOwns Keys in Range
S0100901 to 999, then 0 to 100
S1300101 to 300
S2500301 to 500
S3700501 to 700
S4900701 to 900

5. Looking Up a Key

To find which server owns a key, we hash the key, place it on the ring, and move clockwise until we reach the first server.

For example, if a profile key lands at 250, the next server is S1 at 300, so S1 owns it. A key at 620 goes to S3 at 700. If a key lands at 950, we pass the end of the ring, wrap back to 0, and reach S0 at 100.

Scroll
Key PositionFirst Server ClockwiseOwner
250S1 at 300S1
620S3 at 700S3
950Wraps past 999 to S0 at 100S0
700Lands exactly on S3S3

If a key lands exactly on a server's position, that server owns it.

6. Adding a Server

Now suppose S5 joins the ring at position 400.

Before that, S2 at 500 owned the keys between 300 and 500. After S5 joins, it takes over the range from 300 to 400.

So a key at 350 now reaches S5 first and moves from S2 to S5, while a key at 450 still goes to S2. Keys in the rest of the ring stay exactly where they were.

Scroll
Key PositionOwner Before S5Owner After S5Moved?
250S1S1No
350S2S5Yes
450S2S2No
620S3S3No
950S0S0No

That is the key advantage over modulo hashing. Adding a server changes only a small part of the mapping instead of reshuffling keys across the entire cluster.

7. Removing a Server

Next, suppose S4 at 900 leaves the ring. It previously owned the keys between 700 and 900.

Once S4 is gone, those keys continue clockwise, wrap around the ring, and reach S0 at 100.

So a key at 850 now moves to S0. But a key at 620 still belongs to S3, and a key at 250 still belongs to S1.

Scroll
Key PositionOwner BeforeOwner After S4 LeavesMoved?
250S1S1No
620S3S3No
850S4S0Yes
950S0S0No

Only the range owned by the failed server needs to move.

8. How Many Keys Move

With evenly distributed keys and balanced server positions, adding one new server to a cluster of N servers moves roughly 1 / (N + 1) of the keys. So when we grow from five servers to six, we would expect about one sixth of the keys to move.

That is only an average, though. The exact amount depends on where the servers land on the ring. In our simple example, S5 took 100 out of 1,000 positions, so only about 10% of the keys moved.

Also, consistent hashing only changes ownership. It does not move the data itself.

For a cache, that is usually fine, because missing entries can be loaded again on demand. For persistent storage, we still need a separate process to safely transfer the affected data to its new owner.

9. The Problem of Uneven Ranges

In our example, the servers were evenly spaced around the ring. In practice, hashing real server names usually creates uneven gaps.

One server may own a small range, while another ends up owning a much larger one and therefore handles more keys, memory, and traffic.

In this ring, S2 owns almost half the positions while S1 owns one twentieth. With evenly distributed keys, S2 handles roughly ten times the traffic of S1.

There is another issue when a server fails. With only one position per server, its entire range moves to the next server clockwise. That neighbor can suddenly take on a large amount of extra load. If S2 above fails, S3 goes from owning 80 positions to owning 570.

A common solution is to give each physical server multiple positions on the ring instead of just one.

10. Virtual Nodes

These extra positions are called virtual nodes, or vnodes.

For example, instead of hashing only S1, we might hash identifiers like S1-0, S1-1, and S1-2. Each one gets its own position on the ring, but they all point back to the same physical server. We do the same for every server, spreading their virtual nodes across the ring.

The same positions, grouped by physical server, look like this.

Scroll
Physical ServerVirtual NodesRing Positions
S1S1-0, S1-1, S1-210, 70, 120
S2S2-0, S2-1, S2-250, 80, 160
S3S3-0, S3-1, S3-230, 90, 140

The lookup rule stays exactly the same. We move clockwise to the first position, then route the request to the physical server behind it.

A key still has only one primary owner. Virtual nodes simply give each server several smaller ranges instead of one large range. They do not create extra machines or duplicate the data.

Why Virtual Nodes Help

Using many smaller ranges usually balances keys more evenly than giving each server just one position. The large and small gaps from hashing individual server names average out across dozens of positions.

If a physical server leaves, its ranges are also spread across multiple neighbors instead of being handed to a single server.

Virtual nodes also let us account for different server capacities. A more powerful server can be assigned more virtual nodes, giving it a larger expected share of the keys.

Scroll
BenefitOne Position per ServerMany Virtual Nodes per Server
BalanceRanges vary widely in sizeMany small ranges average out
Failure handlingWhole range lands on one neighborRanges spread across several neighbors
Capacity weightingNot possibleBigger servers get more virtual nodes
CostTiny routing tableMore metadata, more work on membership changes

But there is a trade-off. More virtual nodes mean more routing metadata and more work when membership changes. So in practice, we choose enough virtual nodes to get good balance and then measure the actual distribution. Tens to a few hundred virtual nodes per server is a common starting point.

Virtual nodes improve balance statistically, but they do not guarantee perfectly equal load across every server.

11. Implementing the Ring

In code, we do not need to represent the ring as an actual circle. We can store the virtual node positions in a sorted array and keep a mapping from each position to its physical server.

To route a key, we hash it and use binary search to find the first position greater than or equal to that hash. If there is no such position, we wrap around and use the first entry in the array. If the ring is empty, there is no server to route to, so we return an error.

We also need to handle hash collisions between virtual nodes deterministically, so that one server does not accidentally overwrite another. With a large hash space, collisions are rare, but a production ring should still pick a fixed rule, such as keeping the virtual node whose server id sorts first. The implementation below keeps things short and leaves that rule to the map it uses.

For production code, use a fast, stable, well-distributed non-cryptographic hash such as MurmurHash, xxHash, or CityHash. The examples use FNV-1a because it is short and gives the same output across languages.

Lookup is O(log V), where V is the number of virtual nodes. Adding or removing a physical server is O(R log V), where R is the number of virtual nodes assigned to that server.

One more requirement follows from this design. Every client that routes keys must use the same server list and the same hash function. If two clients see different membership, they can send the same key to different owners, and the cache stops behaving consistently. Systems usually publish ring updates from one place, such as a configuration service, so all clients converge on the same view.

12. Replication on the Ring

So far, each key has had a single owner. If we want to replicate the data for higher availability, we can keep three copies of each key. We choose the primary owner, then continue clockwise to find two more distinct physical servers.

The word distinct is important. The next virtual node might belong to the same machine, and storing another copy there would not help if that machine fails.

Scroll
Key at 75, Walking ClockwiseVirtual NodePhysical ServerAction
First positionS2-1 at 80S2Primary
Second positionS3-1 at 90S3Replica 1
Third positionS1-2 at 120S1Replica 2
If the third had been S2-2S2-2S2Skip, already used

In practice, we may also spread replicas across different racks or availability zones, so that one rack losing power does not take out all three copies.

The ring only helps decide where replicas should live. Other parts of the system still need to copy the data, decide when a write is considered successful, and repair missing replicas.

13. When to Use Consistent Hashing

Consistent hashing works well for caches, storage shards, and other systems where each key should usually go to the same server. If any server can handle any request, a normal load balancer is often simpler.

There are other ways to keep routing stable too. Rendezvous hashing scores each server for a key and picks the highest one. Another common approach is to hash keys into a fixed number of logical partitions, then assign those partitions to servers. Modulo hashing works fine in that case, because the partition count does not change.

Scroll
TechniqueHow It Picks an OwnerGood Fit
Consistent hashing ringFirst server clockwise from the key's positionCaches, storage shards, stateful routing
Rendezvous hashingScore every server for the key, pick the highestSmall to medium server sets, no ring to maintain
Fixed partitionsKey maps to a partition, a table maps partitions to serversLogs, queues, databases that move whole partitions
Plain load balancingAny healthy serverStateless services where any server can handle any request

Summary

Consistent hashing gives keys stable owners in a cluster whose servers change. Plain modulo hashing moves most keys whenever the live server count changes, which is expensive for caches and stateful stores.

Consistent hashing maps both keys and servers onto a fixed ring and assigns each key to the next server clockwise from its position. Adding a server moves only the keys in that server's new range. Removing a server moves only the keys it owned. With N balanced servers, adding one moves about 1 / (N + 1) of the keys.

Virtual nodes give each physical server many positions on the ring. That evens out the ranges, spreads a failed server's load across several neighbors, and lets bigger servers take a larger share.

Consistent hashing decides placement. It does not move data, copy replicas, or decide when a write is durable. Those are separate parts of the system.

As an exercise, draw a small ring yourself with four or five servers and a handful of keys. Add a server and see which keys move. Then remove one and follow those keys to their new owners.

Quiz

Consistent Hashing Quiz

10 quizzes