Consistent Hashing

⭐ Interview Importance: HIGH
⏱️ Revision Time: 3 min

Concept

In standard Hash-Based Sharding, you route data using a modulo operation based on the total number of servers.
server_index = hash(user_id) % 4_servers

This works perfectly to distribute load evenly across Server 0, 1, 2, and 3.
But what happens when your company grows, and you need to add a 5th server?

The formula changes to: hash(user_id) % 5_servers.
Because the math fundamentally changed, a user who previously routed to Server 2 might now route to Server 4. You literally have to physically migrate 80% of your entire database to new servers to satisfy the new math. Your database goes offline for 3 days.

Consistent Hashing is a brilliant mathematical algorithm invented by MIT in 1997 to solve this exact problem. It allows you to add or remove servers with virtually zero data movement.

How It Works (The Hash Ring)

Instead of a standard modulo array [0, 1, 2, 3], Consistent Hashing envisions the hash space as a giant, connected Circle (The Hash Ring).
Imagine the circle goes from 0 to 359 degrees.

Step 1: Place the Servers

You run the IP addresses of your 4 servers through a hash function and place them randomly on the ring.

  • Server A is at 0 degrees.
  • Server B is at 90 degrees.
  • Server C is at 180 degrees.
  • Server D is at 270 degrees.

Step 2: Place the Data

When a user signs up, you hash their user_id and place them on the ring (e.g., User 1 lands at 45 degrees).
The Rule: To find the data, you start at the user’s position (45) and walk clockwise around the circle until you hit the very first Server. User 1 hits Server B (90 degrees).

Step 3: Add a Server (The Magic)

You buy a new Server E. You hash it, and it lands at 135 degrees (right between Server B and Server C).
Who is affected? Only the users sitting between 90 and 135 degrees. Instead of walking to Server C, they now walk to Server E.
Result: You only have to migrate data off Server C. Servers A, B, and D are completely untouched. Adding a new server only required moving 1/N of the total data, preventing massive database downtime.

Virtual Nodes

There is one flaw with the basic ring. What if the hash function accidentally places Server A, B, and C right next to each other (at 1, 2, and 3 degrees), and leaves a massive empty gap until Server D at 200 degrees? Server D will unfairly absorb 60% of all the traffic.

To fix this, modern systems use Virtual Nodes.
Instead of placing “Server A” on the ring exactly once, the algorithm generates 100 fake variations of it (“Server A_1”, “Server A_2”) and scatters them randomly all over the ring. It does this for every server.
By scattering thousands of Virtual Nodes around the circle, the gaps mathematically average out, guaranteeing perfectly even load distribution regardless of where the servers land.

Interview Questions

Q: Consistent Hashing is most famous for powering which massively scalable NoSQL databases?
A: Cassandra and Amazon DynamoDB. These databases use Consistent Hashing as their absolute foundational architecture, allowing them to add and remove hundreds of nodes to a live cluster continuously without the developers or users ever noticing a blip in performance.

Q: A caching layer (like Memcached or Redis) is distributed across 5 servers using standard modulo hashing (% 5). Server 3 experiences a hardware failure and dies. What catastrophic event occurs?
A: Because the total server count drops from 5 to 4, the modulo math changes for every single user. 80% of the active users are suddenly routed to the wrong cache server, resulting in an 80% Cache Miss Rate simultaneously. All of those users will instantly query the backend relational database, causing a massive Thundering Herd that crashes the primary database. If the developers had used Consistent Hashing for their Memcached routing, the failure of Server 3 would have only caused a 20% cache miss rate (smoothly distributing Server 3’s lost traffic to the other nodes), keeping the database safe.