Today we are going to learn about one of the most genius tricks ever invented in computer science — a trick so clever that every major tech company (Google, Amazon, Netflix, Discord) quietly uses it behind the scenes every single day.
When you open YouTube and a video loads in 2 seconds — that's Consistent Hashing at work.
When Discord servers 150 million users simultaneously — Consistent Hashing is the hero.
When Amazon never loses your cart even during peak sale days — you guessed it. 🛒
🍕 Let's Start With a Pizza Shop Story
Imagine you own a pizza delivery business in a big city. 🏙️
You have 4 delivery boys — Amit, Bunty, Chotu, and Dinesh.
Every time an order comes in, you need to decide: which delivery boy should handle it?
You want to be fair. You also want the same delivery boy to always serve the same area
— so he learns the roads, delivers faster, and customers are happy.
This is exactly the problem that Consistent Hashing solves — but instead of pizza orders, it handles billions of internet requests every day! 🌐
🔢 First, Let's Understand "Normal" Hashing (The Old Way)
Before Consistent Hashing existed, engineers used a simple formula:
Server to use = (request number) ÷ (total number of servers) → take the remainder.
This is called Modulo Hashing. Let's see it in action.
A hash is just a number that represents something.
You take any value (like a user's name "Rahul" or a URL) and run it through a hash function — a math formula — and it spits out a number.
"Rahul" → hash function → 47
"Priya" → hash function → 83
Same input always gives the same output. Magic! ✨
Now suppose we have 3 servers (Server 0, Server 1, Server 2).
User "Rahul" → hash = 47 → 47 ÷ 3 = remainder 2 → Goes to Server 2
User "Priya" → hash = 83 → 83 ÷ 3 = remainder 2 → Goes to Server 2
User "Anil" → hash = 60 → 60 ÷ 3 = remainder 0 → Goes to Server 0
User "Meena" → hash = 91 → 91 ÷ 3 = remainder 1 → Goes to Server 1
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Server 0 │ │ Server 1 │ │ Server 2 │
│ (Anil) │ │ (Meena) │ │ (Rahul) │
└──────────┘ └──────────┘ │ (Priya) │
└──────────┘
This works beautifully — as long as the number of servers never changes.
But what happens when a server crashes? Or when we add a new server?
💥 The Big Problem — When Servers Come and Go
Imagine you are having a great day. Everything is running smoothly.
Then suddenly — Server 1 dies! 💀
Now you have only 2 servers left. The formula changes from ÷ 3 to ÷ 2.
BEFORE (3 servers): AFTER Server 1 dies (2 servers): Rahul → hash 47 → 47÷3=2 → Server 2 Rahul → hash 47 → 47÷2=1 → Server 1 ❗ Priya → hash 83 → 83÷3=2 → Server 2 Priya → hash 83 → 83÷2=1 → Server 1 ❗ Anil → hash 60 → 60÷3=0 → Server 0 Anil → hash 60 → 60÷2=0 → Server 0 ✅ Meena → hash 91 → 91÷3=1 → Server 1 Meena → hash 91 → 91÷2=1 → Server 1 ❗ ⚠️ Almost EVERYONE gets shuffled to a different server! ⚠️ All cached data on previous servers is now USELESS — CACHE MISS for
everyone!
When any server is added or removed, almost ALL requests get remapped to different servers.
This means every user suddenly hits a server that has no cached data for them.
Result: Massive slowdown. Overloaded databases. Angry users. 😤
For a company like Netflix with 1000+ servers, adding just ONE server would cause a cache storm that could crash everything.
💡 Enter Consistent Hashing — The Game Changer
Consistent Hashing was invented to solve exactly this problem.
The brilliant idea? Arrange everything in a circle! 🔵
Instead of a straight line of servers, imagine a giant clock face — or a ring — with numbers going from 0 all the way around back to 0. This is called the Hash Ring.
Place both servers and requests on the same giant circular ring using a hash function — then each request travels clockwise to find the nearest server.
🎡 Step 1 — Build the Ring
Imagine a circle with positions from 0 to 360 (like a clock).
We run each server's name through a hash function to get a position on the ring.
hash("Server-A") = 30 → placed at position 30 on the ring
hash("Server-B") = 120 → placed at position 120 on the ring
hash("Server-C") = 220 → placed at position 220 on the ring
hash("Server-D") = 310 → placed at position 310 on the ring
0 / 360
|
Server-D 310 30 Server-A
\ /
\ /
220 --- HASH RING --- (going around)
/ \
/ \
Server-C 220 120 Server-B
📍 Step 2 — Place Requests on the Ring Too
Now when a user request comes in, we hash it too — and it gets a position on the same ring!
Then the rule is super simple: go clockwise from your position
until you hit the first server. That server handles your request.
hash("Rahul's request") = 60 → position 60 → go clockwise → hits Server-B (120) ✅
hash("Priya's request") = 170 → position 170 → go clockwise → hits Server-C (220) ✅
hash("Anil's request") = 250 → position 250 → go clockwise → hits Server-D (310) ✅
hash("Meena's request") = 340 → position 340 → go clockwise → wraps around →
hits Server-A (30) ✅
0 / 360
|
Server-D 310 30 Server-A ← Meena's request (340) wraps around here
\ /
\ /
250 --- RING --- 60 ← Rahul goes to Server-B
Anil↗ ↖Rahul
/ \
/ \
Server-C Server-B
(170) (60→120)
↑Priya's req hits here
🧪 The Magic Moment — What Happens When a Server Crashes?
This is where Consistent Hashing becomes truly beautiful.
Let's say Server-B crashes. What happens?
BEFORE Server-B crashes: AFTER Server-B crashes: Rahul (pos 60) → Server-B (120) ✅ Rahul (pos 60) → skip B → Server-C (220) ♻️ Priya (pos 170) → Server-C (220) ✅ Priya (pos 170) → Server-C (220) ✅ (unchanged!) Anil (pos 250) → Server-D (310) ✅ Anil (pos 250) → Server-D (310) ✅ (unchanged!) Meena (pos 340) → Server-A (30) ✅ Meena (pos 340) → Server-A (30) ✅ (unchanged!) ✨ Only Rahul's request moved — because Server-B was HIS server. ✨ Everyone else? Completely unaffected. Zero disruption!
When a server is added or removed, only the requests that were directly assigned to that server need to be moved.
All other requests stay exactly where they were. 🏆
With N servers and K requests, only K/N requests need remapping — compared to modulo hashing which remaps almost all K requests!
🔢 Let's Understand the Math (Simply!)
Let's say we have 100 requests and 4 servers.
Old Modulo Hashing — when 1 server is removed: ┌─────────────────────────────────────────────────────────┐ │ Requests remapped = ~75 out of 100 (75% disruption!) │ └─────────────────────────────────────────────────────────┘ Consistent Hashing — when 1 server is removed: ┌─────────────────────────────────────────────────────────┐ │ Requests remapped = ~25 out of 100 (only 25%!) │ │ Formula: K/N = 100/4 = 25 requests move. That's it! │ └─────────────────────────────────────────────────────────┘ At Netflix scale (1 billion requests, 1000 servers): Modulo: ~999 million requests disrupted 😱 Consistent: ~1 million requests disrupted 🎉
⚠️ One Problem — Uneven Distribution (Hot Spots!)
The simple ring works great — but there is one sneaky problem. 🕵️
Servers might not be placed evenly on the ring. By random chance,
one server might cover a huge arc while another covers a tiny arc.
BAD distribution (by bad luck): 0 -------- Server-A (50) ---- Server-B (60) ---- Server-C (350) ---- 360 Arc covered: Server-A: 50 positions (positions 350 → 50) ← small slice 🍕 Server-B: 10 positions (positions 50 → 60) ← tiny slice! 😰 Server-C: 290 positions (positions 60 → 350) ← MASSIVE slice 🏋️ Result: Server-C gets ~80% of all traffic. It's OVERLOADED. HOT SPOT! 🔥
When servers land unevenly on the ring, some servers get way more traffic than others.
This is called a hot spot — one server is on fire 🔥 while others are sitting idle.
This defeats the entire purpose of having multiple servers!
🌟 The Solution — Virtual Nodes (VNodes)
Here's where engineers get really clever. Instead of placing each server once on the ring, we place every server many times — at different positions! 🗺️
Each extra copy of a server on the ring is called a Virtual Node (or VNode).
WITHOUT VNodes (1 position each):
Ring: ----[A]----------[B]--[C]------------------------------------- (very uneven!)
WITH VNodes (3 positions each):
Ring: --[A1]--[B1]--[C1]--[A2]--[B2]--[C2]--[A3]--[B3]--[C3]-- (nicely spread!)
How we generate VNode positions:
hash("Server-A-replica-1") = position 30
hash("Server-A-replica-2") = position 145
hash("Server-A-replica-3") = position 270
→ All three map back to the SAME physical Server-A machine
More VNodes = Better distribution = Less hot spots! 🎯
In production systems today, each physical server typically gets 100 to 200 virtual nodes on the ring.
The more VNodes, the more even the distribution — but too many uses more memory.
Most modern databases (Cassandra, DynamoDB) default to 150 VNodes per server.
💻 Code Time! Let's Build Consistent Hashing From Scratch
We are going to build a Consistent Hashing Ring in Python from scratch.
Think of it like building a merry-go-round 🎠 where servers sit at specific spots, and every incoming request spins around until it finds the nearest server.
The code does 3 things:
1️⃣ Lets you add servers to the ring (with virtual nodes)
2️⃣ Lets you remove servers from the ring
3️⃣ Given a request, finds the correct server for it
import hashlib
import bisect
class ConsistentHashRing:
def __init__(self, virtual_nodes=150):
"""
virtual_nodes: how many copies of each server we put on the ring
150 is the industry standard (used by Apache Cassandra)
"""
self.virtual_nodes = virtual_nodes
self.ring = {} # maps hash_position → server_name
self.sorted_keys = [] # sorted list of all positions on the ring
def _hash(self, key):
"""
Takes any string (server name or request key)
and returns a number between 0 and 2^32.
Think of this as: "where on the ring does this go?"
"""
return int(hashlib.md5(key.encode()).hexdigest(), 16) % (2**32)
def add_server(self, server_name):
"""
Adds a server to the ring — at MULTIPLE positions (virtual nodes).
Like placing 150 little flags for this server all around the ring.
"""
for i in range(self.virtual_nodes):
virtual_key = f"{server_name}-replica-{i}"
hash_position = self._hash(virtual_key)
self.ring[hash_position] = server_name
bisect.insort(self.sorted_keys, hash_position)
print(f"✅ Added {server_name} at {self.virtual_nodes} positions
on the ring")
def remove_server(self, server_name):
"""
Removes a server from the ring — takes down ALL its virtual nodes.
Only the requests that were going to this server will be reassigned.
"""
for i in range(self.virtual_nodes):
virtual_key = f"{server_name}-replica-{i}"
hash_position = self._hash(virtual_key)
del self.ring[hash_position]
self.sorted_keys.remove(hash_position)
print(f"🗑️ Removed {server_name} from the ring")
def get_server(self, request_key):
"""
Given a request (like a user ID or URL),
finds which server should handle it.
Rule: go CLOCKWISE from the request's position → hit the nearest
server.
"""
if not self.ring:
return None
hash_position = self._hash(request_key)
# bisect_right finds the next server clockwise
index = bisect.bisect_right(self.sorted_keys, hash_position)
# If we go past the end of the ring, wrap around to start
(it's a circle!)
if index == len(self.sorted_keys):
index = 0
server_position = self.sorted_keys[index]
return self.ring[server_position]
Now we will test our ring to see it in action! 🎉
We will add 3 servers, send 6 user requests, then remove 1 server and watch how only a few requests change — the rest stay put!
# ---- Testing our Consistent Hash Ring ----
ring = ConsistentHashRing(virtual_nodes=150)
# Step 1: Add 3 servers
ring.add_server("Server-Mumbai")
ring.add_server("Server-Delhi")
ring.add_server("Server-Bangalore")
# Step 2: See which server handles each user
users = ["user:Rahul", "user:Priya", "user:Anil", "user:Meena", "user:Kiran", "user:Zara"]
print("\n--- BEFORE removing a server ---")
mapping_before = {}
for user in users:
server = ring.get_server(user)
mapping_before[user] = server
print(f" {user:20s} → {server}")
# Step 3: Remove one server (simulate crash or maintenance)
print("\n")
ring.remove_server("Server-Delhi")
# Step 4: See which users were affected
print("\n--- AFTER removing Server-Delhi ---")
moved = 0
for user in users:
server = ring.get_server(user)
changed = "♻️ MOVED" if server != mapping_before[user] else "✅ same"
print(f" {user:20s} → {server:25s} {changed}")
if server != mapping_before[user]:
moved += 1
print(f"\n Total moved: {moved}/{len(users)} users ({moved/len(users)*100:.0f}%)")
print(f" Expected: ~{len(users)//3}/{len(users)} users ({100//3}%)")
Sample Output:
✅ Added Server-Mumbai at 150 positions on the ring
✅ Added Server-Delhi at 150 positions on the ring
✅ Added Server-Bangalore at 150 positions on the ring
--- BEFORE removing a server ---
user:Rahul → Server-Delhi
user:Priya → Server-Mumbai
user:Anil → Server-Bangalore
user:Meena → Server-Delhi
user:Kiran → Server-Bangalore
user:Zara → Server-Mumbai
🗑️ Removed Server-Delhi from the ring
--- AFTER removing Server-Delhi ---
user:Rahul → Server-Bangalore ♻️ MOVED
user:Priya → Server-Mumbai ✅ same
user:Anil → Server-Bangalore ✅ same
user:Meena → Server-Bangalore ♻️ MOVED
user:Kiran → Server-Bangalore ✅ same
user:Zara → Server-Mumbai ✅ same
Total moved: 2/6 users (33%)
Expected: ~2/6 users (33%) ← Theory matches reality! 🎯
🌍 Real-World Consistent Hashing — How Companies Use It
📦 1. Apache Cassandra (Used by Instagram, Netflix)
Cassandra is a distributed database that stores data across many machines.
It uses Consistent Hashing to decide which machine stores which rows of data.
Each row has a partition key (like a user ID). Cassandra hashes this key and finds the correct machine on the ring. When Instagram adds a new database node, only the rows from the adjacent node need to be moved — zero downtime! ⚡
Instagram User Table in Cassandra: Partition Key (user_id) → hash → position on ring → stored on node user_id: 1001 → hash: 45 → stored on Node-A (position 50) user_id: 1002 → hash: 130 → stored on Node-B (position 140) user_id: 1003 → hash: 280 → stored on Node-C (position 300) Add new Node-D at position 200: → Only rows between position 140 and 200 move from Node-C to Node-D → All other data stays put! 🏆
⚡ 2. Amazon DynamoDB (Used by Airbnb, Lyft)
DynamoDB is Amazon's lightning-fast database service.
Behind the scenes, it uses Consistent Hashing with hundreds of virtual nodes
to spread your data perfectly across its global infrastructure. 🌐
When you book a ride on Lyft, your trip data goes through a hash function, lands on a specific DynamoDB node, and you never wait more than a millisecond. ⏱️
🗄️ 3. Redis Cluster (Used by Twitter, GitHub)
Redis is an ultra-fast in-memory cache. When scaled to a cluster, it uses a concept called hash slots — basically Consistent Hashing with 16,384 slots on the ring — to distribute keys across nodes.
Redis Cluster Hash Slots:
Total slots = 16,384
Node 1 (Master A): handles slots 0 – 5460
Node 2 (Master B): handles slots 5461 – 10922
Node 3 (Master C): handles slots 10923– 16383
GET "tweet:rahul:12345"
→ CRC16("tweet:rahul:12345") % 16384 = 7291
→ 7291 falls in range 5461–10922 → goes to Node 2 ✅
Adding Node 4? Just redistribute some slots. Smooth. Zero drama. 😎
🌐 4. Content Delivery Networks — CDN (Used by Cloudflare, Akamai)
When you watch a YouTube video, it doesn't come from Google's headquarters in California.
It comes from a CDN server near you — maybe in Mumbai or Hyderabad! 🏙️
CDN systems use Consistent Hashing to route your video request to the
nearest edge server that has your video cached.
If that edge server is down, only its users are rerouted — others are unaffected.
🔁 Replication — Don't Put All Eggs in One Basket
What if the server holding your data crashes permanently? 😱
That's why smart systems don't just put data on one server —
they replicate it to the next N servers clockwise on the ring!
Replication Factor = 3 (Cassandra default)
user_id: 1001 → hash position 45 → primary node: Server-A (pos 50)
→ replica 1: Server-B (pos 120) ← next clockwise
→ replica 2: Server-C (pos 220) ← next clockwise
again
Even if Server-A AND Server-B both crash,
your data is still safe on Server-C! 💪
Ring view:
───[Server-A(50)]───[Server-B(120)]───[Server-C(220)]───[Server-D(310)]───
PRIMARY REPLICA-1 REPLICA-2
←── data for user 1001 stored on these 3 ───────────────────────────────►
🏦 Banking systems: Replication Factor = 5 (ultra-safe)
📱 Social media (Instagram, Twitter): Replication Factor = 3 (balanced)
🎮 Gaming leaderboards: Replication Factor = 2 (speed over safety)
Cassandra lets you configure this per table! So you decide the trade-off.
📐 The Full System Design — Putting It All Together
Now let's see how a real production system uses Consistent Hashing end to end.
We will use the example of a distributed caching system like Memcached or Redis.
Full Request Flow:
[User Request: "Get profile for user_id=9876"]
│
▼
[Load Balancer / API Gateway]
→ Receives the request, passes it to the Cache Coordinator
│
▼
[Cache Coordinator — Consistent Hash Ring]
→ hash("user:9876") = 178
→ Looks up ring → next clockwise server = Cache-Node-C (position 220)
│
▼
[Cache-Node-C]
→ Is user 9876's profile stored here? (Cache HIT ✅ or MISS ❌)
│ │
HIT ✅ MISS ❌
│ │
▼ ▼
Return cached data Query the Database
in < 1ms ⚡ → Store result in Cache-Node-C
→ Return data to user
→ Next time: Cache HIT! ✅
Key insight: user 9876 ALWAYS goes to Cache-Node-C (unless C is removed).
So the cached data is always in the right place! No wasted cache lookups. 🎯
🆚 Consistent Hashing vs Alternatives — Comparison
┌────────────────────┬───────────┬──────────────┬──────────────────────────┐ │ Method │ Add Server│ Remove Server│ Who Uses It │ ├────────────────────┼───────────┼──────────────┼──────────────────────────┤ │ Modulo Hashing │ Remaps │ Remaps ~ALL │ Old systems, simple apps │ │ │ ~ALL keys │ keys 😱 │ │ ├────────────────────┼───────────┼──────────────┼──────────────────────────┤ │ Consistent Hashing │ Remaps │ Remaps K/N │ Cassandra, DynamoDB, │ │ (Ring) │ K/N keys │ keys only ✅ │ Redis Cluster, Memcached │ ├────────────────────┼───────────┼──────────────┼──────────────────────────┤ │ Rendezvous Hashing │ Remaps │ Remaps K/N │ CDNs, Nginx, HAProxy │ │ (HRW) │ K/N keys │ keys only ✅ │ (simpler implementation) │ ├────────────────────┼───────────┼──────────────┼──────────────────────────┤ │ Jump Consistent │ Remaps │ N/A │ Google Spanner, FoundDB │ │ Hashing │ K/N keys │ (add only) │ (fastest lookup O(ln N)) │ └────────────────────┴───────────┴──────────────┴──────────────────────────┘
Modern systems like Apache Kafka 4.x and FoundationDB are moving toward Range-Based Partitioning combined with Consistent Hashing for even better control over data locality.
But the classic Hash Ring remains the most widely deployed technique by far — powering systems that serve trillions of requests daily.
🧠 Advanced: Bounded Loads — The Enhancement
Even with VNodes, sometimes one server gets slightly more traffic due to hot keys (super popular users like a celebrity's account). 📸
Google researchers published a technique in 2017 called Consistent Hashing with Bounded Loads. The idea: each server has a maximum load limit. If the primary server is too full, the request overflows to the next server on the ring.
Standard Consistent Hashing: Viral tweet by @celebrity → hash → Server-A → Server-A gets 1 million hits 🔥 Bounded Load Consistent Hashing: Viral tweet by @celebrity → hash → Server-A → Server-A is above 110% capacity? → Overflow to Server-B (next clockwise) → Server-B also full? → Overflow to Server-C → Load is now spread! Capacity limit formula: each server accepts at most (1 + ε) × average_load Where ε (epsilon) = your tolerance level (typically 0.1 to 0.25 in production)
→ Google Cloud Load Balancer — uses it for traffic distribution
→ Meta's Proxygen — HTTP proxy handling billions of social media requests
→ Uber's Ringpop — coordinates microservices across their fleet
This is the cutting-edge version of Consistent Hashing that you should know for senior interviews!
🎯 System Design Interview — How to Use This Knowledge
When an interviewer asks: "Design a distributed cache like Memcached" — here is your winning answer structure:
Step 1: "I'll use Consistent Hashing to distribute keys across cache nodes."
→ Explain the ring concept
→ Mention why modulo hashing fails at scale
Step 2: "To prevent hot spots, I'll use 150 virtual nodes per server."
→ Mention Cassandra uses the same approach
→ Shows you know production systems
Step 3: "For fault tolerance, I'll replicate each key to the next 2 clockwise nodes."
→ Replication factor of 3
→ Data survives 2 node failures
Step 4: "To handle viral content (hot keys), I'll add bounded load limits."
→ Shows you know the state-of-the-art
→ Interviewers will be VERY impressed 🏆
Bonus: "I'd use a gossip protocol like Cassandra's to let nodes discover
each other's ring positions without a central coordinator."
→ This shows real distributed systems depth!
🚀 Quick Recap
🔴 Problem: Servers crash. New servers are added. Modulo hashing remaps everything. ✅ Solution: Consistent Hashing — a ring where only K/N requests move on changes. 🔴 Problem: Servers land unevenly on the ring → hot spots. ✅ Solution: Virtual Nodes (VNodes) — 150 copies per server, perfectly spread. 🔴 Problem: Servers crash, data is lost. ✅ Solution: Replication — store data on next N clockwise servers (usually N=3). 🔴 Problem: Viral content / hot keys overload one server. ✅ Solution: Bounded Load Hashing — overflow to next server when capacity exceeded. 🏆 Who uses it? Cassandra, DynamoDB, Redis, Nginx, Cloudflare, Discord, Netflix, Google. 🎓 You are now at HERO level. Go ace that interview! 💪
Ring → Hash → Clockwise → VNodes → Replicate 🎯
Say these 5 words before your next system design interview. You now understand the same technology that powers the apps used by 4 billion people every day. That's pretty amazing for one blog post! 🌟
Happy Hashing! 🔗✨
Comments
Post a Comment