Skip to main content

Consistent Hashing Explained: How Distributed Systems Scale Efficiently

Calculating read time…

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.

💡 Why Should You Care?
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.

💡 What is a Hash?
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!
❌ The Modulo Problem — Why It Fails at Scale:
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.

✅ The Core Idea in One Sentence:
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!
✅ The Golden Rule of Consistent Hashing:
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! 🔥
❌ The Hot Spot Problem:
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! 🎯
✅ VNode Rule of Thumb ( Best Practice):
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

🗒️ What the code below does — Read this first!
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]
🗒️ What the code below does — Read this first!
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 ───────────────────────────────►
✅ Industry Standard Replication Factors :
🏦 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)) │
  └────────────────────┴───────────┴──────────────┴──────────────────────────┘
💡 Trend:
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)
✅ Who Uses Bounded Load Hashing ?
→ 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! 💪
✅ The 5 Words to Remember:
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