Consistent Hashing Explained in Simple Terms
Modern applications often store data across multiple servers. For example, an e-commerce application may use several cache servers to store product details, user sessions, and frequently accessed data.
But this creates an important question:
How do we decide which server should store a particular key?
This is where hashing is used. However, traditional hashing creates a serious problem whenever servers are added or removed. Consistent Hashing solves that problem.
Let us understand it step by step.
1. How Normal Hashing Works
Suppose we have three cache servers:
Server 0
Server 1
Server 2
To find the server for a key, we use:
serverIndex = hash(key) % numberOfServers
Suppose:
hash("user123") = 17
numberOfServers = 3
Therefore:
17 % 3 = 2
The data for user123 will be stored on Server 2.
Here are some more examples:
| Key | Hash value | Calculation | Assigned server |
|---|---|---|---|
| user1 | 10 | 10 % 3 | Server 1 |
| user2 | 14 | 14 % 3 | Server 2 |
| user3 | 21 | 21 % 3 | Server 0 |
| user4 | 25 | 25 % 3 | Server 1 |
As long as the number of servers remains three, the same key will always be mapped to the same server.
2. The Problem With Normal Hashing
The formula depends on the total number of servers:
hash(key) % numberOfServers
Suppose traffic increases and we add a fourth server.
The formula changes from:
hash(key) % 3
to:
hash(key) % 4
Now observe what happens:
| Key | Hash | With 3 servers | With 4 servers |
|---|---|---|---|
| user1 | 10 | Server 1 | Server 2 |
| user2 | 14 | Server 2 | Server 2 |
| user3 | 21 | Server 0 | Server 1 |
| user4 | 25 | Server 1 | Server 1 |
| user5 | 29 | Server 2 | Server 1 |
Many keys are now mapped to different servers.
For example, user1 was previously stored on Server 1, but after adding a new server, the application searches for it on Server 2.
Data is stored on Server 1
Application searches on Server 2
The application will not find the data, resulting in a cache miss.
If millions of keys are remapped, it can cause:
A sudden increase in cache misses
Heavy load on the database
Increased response time
Expensive movement of data between servers
Possible database or application failure
This problem is called massive remapping or massive rehashing.
3. What Is Consistent Hashing?
Consistent Hashing is a technique that distributes data across multiple servers while minimizing the number of keys that must move when a server is added or removed.
Instead of using:
hash(key) % numberOfServers
Consistent Hashing places both servers and keys inside a fixed hash space.
The basic process is:
Create a circular hash space.
Place servers on the circle using their hash values.
Place keys on the same circle.
Move clockwise from the key.
Assign the key to the first server encountered.
This circular structure is called a hash ring.
4. Creating the Hash Ring
Suppose our hash function generates values between 0 and 99:
0, 1, 2, 3, ... 98, 99
We connect the end of the range back to the beginning:
0 → 1 → 2 → ... → 98 → 99 → 0
It works like a clock. After 12, we return to 1. Similarly, after 99, the hash ring returns to 0.
The actual hash range in real systems can be much larger, such as:
0 to 2³² - 1
The small range of 0 to 99 is only being used to make the concept easier to understand.
5. Placing Servers on the Hash Ring
Suppose we have three servers:
Server A
Server B
Server C
We calculate the hash of each server’s identifier, such as its name or IP address.
Assume that we get:
hash("Server A") = 20
hash("Server B") = 50
hash("Server C") = 80
Therefore, the servers are placed at the following positions:
| Server | Hash-ring position |
|---|---|
| Server A | 20 |
| Server B | 50 |
| Server C | 80 |
Their order on the ring is:
0 → A(20) → B(50) → C(80) → 99 → 0
6. Placing Keys on the Hash Ring
Keys are processed using the same hash function.
Suppose:
hash("user1") = 10
hash("user2") = 30
hash("user3") = 65
hash("user4") = 90
The assignment rule is:
Start from the key’s position and move clockwise. Assign the key to the first server found.
Assigning user1
hash("user1") = 10
The first server after position 10 is Server A at position 20.
user1 → Server A
Assigning user2
hash("user2") = 30
The first server after position 30 is Server B at position 50.
user2 → Server B
Assigning user3
hash("user3") = 65
The first server after position 65 is Server C at position 80.
user3 → Server C
Assigning user4
hash("user4") = 90
There is no server between positions 90 and 99. Therefore, we wrap around to position 0.
The first server after wrapping around is Server A at position 20.
user4 → Server A
The final distribution is:
| Key | Hash position | Assigned server |
|---|---|---|
| user1 | 10 | Server A |
| user2 | 30 | Server B |
| user3 | 65 | Server C |
| user4 | 90 | Server A |
7. What Happens When a New Server Is Added?
Suppose we add Server D.
hash("Server D") = 40
The new order becomes:
0 → A(20) → D(40) → B(50) → C(80) → 99 → 0
Let us check the keys again.
user1
Position 10 → Server A at 20
No change occurs.
user2
Position 30 → Server D at 40
Previously, user2 was assigned to Server B. It will now move to Server D.
user3
Position 65 → Server C at 80
No change occurs.
user4
Position 90 → wrap around → Server A at 20
No change occurs.
The result is:
| Key | Before adding D | After adding D |
|---|---|---|
| user1 | Server A | Server A |
| user2 | Server B | Server D |
| user3 | Server C | Server C |
| user4 | Server A | Server A |
Only user2 moved to the new server. All other mappings remained unchanged.
This is the main advantage of Consistent Hashing:
Adding a server affects only a small portion of the keys instead of remapping almost every key.
8. What Happens When a Server Is Removed?
Suppose Server B fails and is removed from the ring.
Before removal:
A(20) → D(40) → B(50) → C(80)
After removal:
A(20) → D(40) → C(80)
The keys that previously belonged to Server B will move clockwise to Server C.
The keys assigned to Servers A and D will remain unchanged.
Therefore:
When a server is removed, only the keys belonging to that server need to move.
This is much better than normal hashing, where changing the server count can change the mapping of almost every key.
9. Why Is the Structure Circular?
Suppose a key is placed at position 90, but the last server is at position 80.
If the structure were a straight line, there would be no server after position 90.
The circular structure solves this problem:
90 → 99 → 0 → first available server
Therefore, every key will always find a server.
10. How Is Data Found Again?
Suppose:
hash("user2") = 30
When writing the data:
Calculate the hash of
user2.Find position 30.
Move clockwise.
Find the next server.
Store the data on that server.
When reading the data, the application performs exactly the same process:
Calculate the same hash.
Reach the same position.
Move in the same clockwise direction.
Find the same server.
As long as the ring has not changed, the key will always reach the same server.
11. The Uneven Distribution Problem
Placing every physical server only once on the ring can result in uneven data distribution.
Suppose the servers are placed at:
Server A = 10
Server B = 20
Server C = 90
The distance between Server B and Server C is much larger than the other ranges. Therefore, Server C may receive significantly more keys.
The result can be:
Server A → Low load
Server B → Low load
Server C → Very high load
Consistent Hashing solves this problem using virtual nodes.
12. What Are Virtual Nodes?
Instead of placing a physical server at only one position, we place it at multiple positions on the ring.
For example:
Server A:
A1 = 10
A2 = 45
A3 = 80
Server B:
B1 = 20
B2 = 55
B3 = 90
Server C:
C1 = 5
C2 = 35
C3 = 70
Here, A1, A2, and A3 are not three different physical servers. They are three virtual positions representing the same Server A.
The ring may look like:
C1 → A1 → B1 → C2 → A2 → B2 → C3 → A3 → B3
Because every physical server appears at multiple locations, data is distributed more evenly.
Benefits of virtual nodes
Virtual nodes provide:
Better load distribution
Easier server addition and removal
Fewer chances of one server becoming overloaded
Flexible distribution based on server capacity
A powerful server can be given more virtual nodes:
Small server → 100 virtual nodes
Medium server → 200 virtual nodes
Large server → 400 virtual nodes
This allows the large server to handle a larger share of the data.
13. Consistent Hashing and Replication
Consistent Hashing determines the primary server for a key. However, if that server fails, the data may become unavailable.
Therefore, distributed systems usually combine Consistent Hashing with replication.
Suppose the replication factor is three:
Replication factor = 3
Starting from the key’s position, the system selects the first three different physical servers in the clockwise direction.
First server → Primary copy
Second server → Replica 1
Third server → Replica 2
For example:
user123 → Server A
Replica → Server B
Replica → Server C
If Server A fails, the data can still be read from Server B or Server C.
14. Normal Hashing vs Consistent Hashing
| Feature | Normal Hashing | Consistent Hashing |
|---|---|---|
| Technique | hash(key) % N |
Hash ring |
| Depends on server count | Yes | No |
| Adding a server | Many keys may move | Only a small range moves |
| Removing a server | Many keys may move | Only removed server’s keys move |
| Load distribution | Simple but unstable | Better with virtual nodes |
| Implementation | Easy | More complex |
| Best suited for | Fixed server count | Dynamic distributed systems |
15. Where Is Consistent Hashing Used?
Consistent Hashing is useful in systems where servers can frequently be added, removed, or replaced.
Common use cases include:
Distributed caching
Database sharding
Content Delivery Networks
Distributed file systems
Session storage
NoSQL databases
Message brokers
Distributed load balancing
It is commonly associated with systems such as:
Apache Cassandra
Amazon Dynamo-style databases
Riak
Distributed Memcached clients
Large Content Delivery Networks
16. Simple Real-Life Analogy
Imagine a circular road with three warehouses:
Warehouse A = Position 20
Warehouse B = Position 50
Warehouse C = Position 80
Customers are also located at different positions on the road.
The rule is:
A customer moves clockwise and uses the first warehouse encountered.
A customer at position 30 will use Warehouse B at position 50.
Now suppose a new Warehouse D opens at position 40. The customer at position 30 will start using Warehouse D.
However, customers in other parts of the road will continue using their existing warehouses.
The new warehouse affects only nearby customers, not every customer on the road.
That is exactly how Consistent Hashing works.
17. Final Summary
Normal hashing uses:
server = hash(key) % numberOfServers
When the number of servers changes, the result of the formula changes for many keys.
Consistent Hashing removes this dependency on the server count.
It follows this process:
Place servers and keys on a hash ring.
Start from the key’s position.
Move clockwise.
Select the first available server.
As a result:
Adding a server moves only a small range of keys.
Removing a server moves only that server’s keys.
Virtual nodes provide better load distribution.
Replication provides fault tolerance.
The most important idea to remember is:
Consistent Hashing does not prevent keys from moving. It minimizes how many keys must move when the server configuration changes
Modern applications often store data across multiple servers. For example, an e-commerce application may use several cache servers to store product details, user sessions, and frequently accessed data.
But this creates an important question:
How do we decide which server should store a particular key?
This is where hashing is used. However, traditional hashing creates a serious problem whenever servers are added or removed. Consistent Hashing solves that problem.
Let us understand it step by step.
1. How Normal Hashing Works
Suppose we have three cache servers:
Server 0
Server 1
Server 2
To find the server for a key, we use:
serverIndex = hash(key) % numberOfServers
Suppose:
hash("user123") = 17
numberOfServers = 3
Therefore:
17 % 3 = 2
The data for user123 will be stored on Server 2.
Here are some more examples:
| Key | Hash value | Calculation | Assigned server |
|---|---|---|---|
| user1 | 10 | 10 % 3 | Server 1 |
| user2 | 14 | 14 % 3 | Server 2 |
| user3 | 21 | 21 % 3 | Server 0 |
| user4 | 25 | 25 % 3 | Server 1 |
As long as the number of servers remains three, the same key will always be mapped to the same server.
2. The Problem With Normal Hashing
The formula depends on the total number of servers:
hash(key) % numberOfServers
Suppose traffic increases and we add a fourth server.
The formula changes from:
hash(key) % 3
to:
hash(key) % 4
Now observe what happens:
| Key | Hash | With 3 servers | With 4 servers |
|---|---|---|---|
| user1 | 10 | Server 1 | Server 2 |
| user2 | 14 | Server 2 | Server 2 |
| user3 | 21 | Server 0 | Server 1 |
| user4 | 25 | Server 1 | Server 1 |
| user5 | 29 | Server 2 | Server 1 |
Many keys are now mapped to different servers.
For example, user1 was previously stored on Server 1, but after adding a new server, the application searches for it on Server 2.
Data is stored on Server 1
Application searches on Server 2
The application will not find the data, resulting in a cache miss.
If millions of keys are remapped, it can cause:
A sudden increase in cache misses
Heavy load on the database
Increased response time
Expensive movement of data between servers
Possible database or application failure
This problem is called massive remapping or massive rehashing.
3. What Is Consistent Hashing?
Consistent Hashing is a technique that distributes data across multiple servers while minimizing the number of keys that must move when a server is added or removed.
Instead of using:
hash(key) % numberOfServers
Consistent Hashing places both servers and keys inside a fixed hash space.
The basic process is:
Create a circular hash space.
Place servers on the circle using their hash values.
Place keys on the same circle.
Move clockwise from the key.
Assign the key to the first server encountered.
This circular structure is called a hash ring.
4. Creating the Hash Ring
Suppose our hash function generates values between 0 and 99:
0, 1, 2, 3, ... 98, 99
We connect the end of the range back to the beginning:
0 → 1 → 2 → ... → 98 → 99 → 0
It works like a clock. After 12, we return to 1. Similarly, after 99, the hash ring returns to 0.
The actual hash range in real systems can be much larger, such as:
0 to 2³² - 1
The small range of 0 to 99 is only being used to make the concept easier to understand.
5. Placing Servers on the Hash Ring
Suppose we have three servers:
Server A
Server B
Server C
We calculate the hash of each server’s identifier, such as its name or IP address.
Assume that we get:
hash("Server A") = 20
hash("Server B") = 50
hash("Server C") = 80
Therefore, the servers are placed at the following positions:
| Server | Hash-ring position |
|---|---|
| Server A | 20 |
| Server B | 50 |
| Server C | 80 |
Their order on the ring is:
0 → A(20) → B(50) → C(80) → 99 → 0
6. Placing Keys on the Hash Ring
Keys are processed using the same hash function.
Suppose:
hash("user1") = 10
hash("user2") = 30
hash("user3") = 65
hash("user4") = 90
The assignment rule is:
Start from the key’s position and move clockwise. Assign the key to the first server found.
Assigning user1
hash("user1") = 10
The first server after position 10 is Server A at position 20.
user1 → Server A
Assigning user2
hash("user2") = 30
The first server after position 30 is Server B at position 50.
user2 → Server B
Assigning user3
hash("user3") = 65
The first server after position 65 is Server C at position 80.
user3 → Server C
Assigning user4
hash("user4") = 90
There is no server between positions 90 and 99. Therefore, we wrap around to position 0.
The first server after wrapping around is Server A at position 20.
user4 → Server A
The final distribution is:
| Key | Hash position | Assigned server |
|---|---|---|
| user1 | 10 | Server A |
| user2 | 30 | Server B |
| user3 | 65 | Server C |
| user4 | 90 | Server A |
7. What Happens When a New Server Is Added?
Suppose we add Server D.
hash("Server D") = 40
The new order becomes:
0 → A(20) → D(40) → B(50) → C(80) → 99 → 0
Let us check the keys again.
user1
Position 10 → Server A at 20
No change occurs.
user2
Position 30 → Server D at 40
Previously, user2 was assigned to Server B. It will now move to Server D.
user3
Position 65 → Server C at 80
No change occurs.
user4
Position 90 → wrap around → Server A at 20
No change occurs.
The result is:
| Key | Before adding D | After adding D |
|---|---|---|
| user1 | Server A | Server A |
| user2 | Server B | Server D |
| user3 | Server C | Server C |
| user4 | Server A | Server A |
Only user2 moved to the new server. All other mappings remained unchanged.
This is the main advantage of Consistent Hashing:
Adding a server affects only a small portion of the keys instead of remapping almost every key.
8. What Happens When a Server Is Removed?
Suppose Server B fails and is removed from the ring.
Before removal:
A(20) → D(40) → B(50) → C(80)
After removal:
A(20) → D(40) → C(80)
The keys that previously belonged to Server B will move clockwise to Server C.
The keys assigned to Servers A and D will remain unchanged.
Therefore:
When a server is removed, only the keys belonging to that server need to move.
This is much better than normal hashing, where changing the server count can change the mapping of almost every key.
9. Why Is the Structure Circular?
Suppose a key is placed at position 90, but the last server is at position 80.
If the structure were a straight line, there would be no server after position 90.
The circular structure solves this problem:
90 → 99 → 0 → first available server
Therefore, every key will always find a server.
10. How Is Data Found Again?
Suppose:
hash("user2") = 30
When writing the data:
Calculate the hash of
user2.Find position 30.
Move clockwise.
Find the next server.
Store the data on that server.
When reading the data, the application performs exactly the same process:
Calculate the same hash.
Reach the same position.
Move in the same clockwise direction.
Find the same server.
As long as the ring has not changed, the key will always reach the same server.
11. The Uneven Distribution Problem
Placing every physical server only once on the ring can result in uneven data distribution.
Suppose the servers are placed at:
Server A = 10
Server B = 20
Server C = 90
The distance between Server B and Server C is much larger than the other ranges. Therefore, Server C may receive significantly more keys.
The result can be:
Server A → Low load
Server B → Low load
Server C → Very high load
Consistent Hashing solves this problem using virtual nodes.
12. What Are Virtual Nodes?
Instead of placing a physical server at only one position, we place it at multiple positions on the ring.
For example:
Server A:
A1 = 10
A2 = 45
A3 = 80
Server B:
B1 = 20
B2 = 55
B3 = 90
Server C:
C1 = 5
C2 = 35
C3 = 70
Here, A1, A2, and A3 are not three different physical servers. They are three virtual positions representing the same Server A.
The ring may look like:
C1 → A1 → B1 → C2 → A2 → B2 → C3 → A3 → B3
Because every physical server appears at multiple locations, data is distributed more evenly.
Benefits of virtual nodes
Virtual nodes provide:
Better load distribution
Easier server addition and removal
Fewer chances of one server becoming overloaded
Flexible distribution based on server capacity
A powerful server can be given more virtual nodes:
Small server → 100 virtual nodes
Medium server → 200 virtual nodes
Large server → 400 virtual nodes
This allows the large server to handle a larger share of the data.
13. Consistent Hashing and Replication
Consistent Hashing determines the primary server for a key. However, if that server fails, the data may become unavailable.
Therefore, distributed systems usually combine Consistent Hashing with replication.
Suppose the replication factor is three:
Replication factor = 3
Starting from the key’s position, the system selects the first three different physical servers in the clockwise direction.
First server → Primary copy
Second server → Replica 1
Third server → Replica 2
For example:
user123 → Server A
Replica → Server B
Replica → Server C
If Server A fails, the data can still be read from Server B or Server C.
14. Normal Hashing vs Consistent Hashing
| Feature | Normal Hashing | Consistent Hashing |
|---|---|---|
| Technique | hash(key) % N |
Hash ring |
| Depends on server count | Yes | No |
| Adding a server | Many keys may move | Only a small range moves |
| Removing a server | Many keys may move | Only removed server’s keys move |
| Load distribution | Simple but unstable | Better with virtual nodes |
| Implementation | Easy | More complex |
| Best suited for | Fixed server count | Dynamic distributed systems |
15. Where Is Consistent Hashing Used?
Consistent Hashing is useful in systems where servers can frequently be added, removed, or replaced.
Common use cases include:
Distributed caching
Database sharding
Content Delivery Networks
Distributed file systems
Session storage
NoSQL databases
Message brokers
Distributed load balancing
It is commonly associated with systems such as:
Apache Cassandra
Amazon Dynamo-style databases
Riak
Distributed Memcached clients
Large Content Delivery Networks
16. Simple Real-Life Analogy
Imagine a circular road with three warehouses:
Warehouse A = Position 20
Warehouse B = Position 50
Warehouse C = Position 80
Customers are also located at different positions on the road.
The rule is:
A customer moves clockwise and uses the first warehouse encountered.
A customer at position 30 will use Warehouse B at position 50.
Now suppose a new Warehouse D opens at position 40. The customer at position 30 will start using Warehouse D.
However, customers in other parts of the road will continue using their existing warehouses.
The new warehouse affects only nearby customers, not every customer on the road.
That is exactly how Consistent Hashing works.
17. Final Summary
Normal hashing uses:
server = hash(key) % numberOfServers
When the number of servers changes, the result of the formula changes for many keys.
Consistent Hashing removes this dependency on the server count.
It follows this process:
Place servers and keys on a hash ring.
Start from the key’s position.
Move clockwise.
Select the first available server.
As a result:
Adding a server moves only a small range of keys.
Removing a server moves only that server’s keys.
Virtual nodes provide better load distribution.
Replication provides fault tolerance.
The most important idea to remember is:
Consistent Hashing does not prevent keys from moving. It minimizes how many keys must move when the server configuration changes.