Skip to main content

Database Sharding in System Design: Strategies, Shard Keys and Scaling Explained

Calculating read time…

Today we learn one of the most important concepts in system design —
used every day by Instagram, WhatsApp, Uber, and thousands of companies.
It is called Database Sharding. 🔀

💡 Our running example throughout this post:

OrderEase — a food delivery app that grew from 1,000 to 500 million users.
The single database is drowning. Queries take 30 seconds. Servers crash daily. 
Sharding saves OrderEase — and keeps it fast forever.


What is Sharding? 

Imagine your school library has ONE librarian and one million books. 📚
Every student in the city comes to this one person for help.
She is always overwhelmed. Students wait hours. Everything is slow.

Now imagine splitting those books across 4 mini-libraries around the city.
Each has its own librarian. Each is fast. Each serves its neighbourhood.
This is sharding — splitting one giant database into smaller, faster pieces.

📊 Watch data being routed to shards in real time

🔀 Shard Router
hash(user_id) % 4 = shard
● routing live data...
↓
uid:4
uid:2
uid:7
uid:3
Shard A
users 0–125M
125M rows
Shard B
users 125–250M
125M rows
Shard C
users 250–375M
125M rows
Shard D
users 375–500M
125M rows

🎬 Pure CSS animation — each coloured packet routes to its correct shard automatically


🗂️ The Four Types of Sharding

Type 1 — Range-Based Sharding

Divide data by a number range.
user_id 1–125M → Shard A. user_id 125M–250M → Shard B. Simple!

  user_id         Shard
  ─────────────   ──────────────────────
  1               Shard A (range 1–125M)
  50,000,000      Shard A
  200,000,000     Shard B (125M–250M)
  350,000,000     Shard C (250M–375M)

  ✅ Simple. Great for range queries.
  ❌ HOTSPOT: new signups fill only the last shard!

Type 2 — Hash-Based Sharding

Run the key through a hash function — a mathematical blender. 🎛️
Result: data spreads perfectly evenly. No hotspots.

⚡ Hash vs Range — Load Distribution

❌ Range Sharding (Launch Day)
Shard A
87%🔴
Shard B
5%
Shard C
5%
Shard D
3%
✅ Hash Sharding (Even spread)
Shard A
25%
Shard B
25%
Shard C
25%
Shard D
25%

Type 3 — Consistent Hashing (The Production Standard)

Regular hash sharding breaks when you add a shard — the formula changes, ALL data must move.
Consistent hashing uses a virtual ring. Adding a shard only moves ~25% of data. 🎯

🔄 Consistent Hash Ring — Rotating Live

A
B
C
D
Hash
Ring

Each shard owns a section of the ring.
Adding Shard E? Only its slice of the ring moves — other shards untouched. ✅
Used by: Amazon DynamoDB · Apache Cassandra · Redis Cluster

Type 4 — Geo / Directory Sharding

Geo Sharding: route users to the shard in their region.
Indian users → Mumbai. Europeans → Frankfurt. US → Virginia.
Data is physically close → blazing-fast reads for every user.

🌍 Geo Sharding — Users Route to Nearest Region

🇮🇳
India Users
↓
Mumbai Shard
~4ms latency ✅
🇪🇺
EU Users
↓
Frankfurt Shard
~6ms latency ✅
🇺🇸
US Users
↓
Virginia Shard
~5ms latency ✅

vs. single DB with no sharding → 180–220ms latency for non-US users 😓


🏗️ The Shard Router — The Traffic Director

Every read/write goes through the Shard Router first.
It answers: "Which shard holds this user's data?"
Think of it as a hotel concierge who knows every guest's room number. 🏨

⚡ Request Flow — From App to Correct Shard

📱 OrderEase App
"Get order for user_id = 8,500,000"
↓
🔀 Shard Router
hash(8500000) % 4 = 2 → Shard C
↓
🗄️ Shard C only
Searches 125M rows — not 500M! → returns in 4ms ✅
↓
✅ Result returned to user
10× faster than a single un-sharded database

💻 Building a Shard Router in Python

📝 What does this code do?

We build a ShardRouter class that takes any user_id and returns the correct database server to connect to.
get_shard() hashes the user_id — same user always maps to same shard.
execute() connects to the right shard and runs your SQL query.
The app never needs to know which shard — the router handles it invisibly.
import hashlib
import psycopg2

class ShardRouter:

    def __init__(self):
        # Each shard = a separate database server in a different data centre
        self.shards = {
            0: {"host": "shard-a.orderease.com", "dbname": "orderease_a"},
            1: {"host": "shard-b.orderease.com", "dbname": "orderease_b"},
            2: {"host": "shard-c.orderease.com", "dbname": "orderease_c"},
            3: {"host": "shard-d.orderease.com", "dbname": "orderease_d"},
        }
        self.num_shards = len(self.shards)

    def get_shard(self, user_id: int) -> int:
        """Same user_id ALWAYS maps to the same shard. Consistent. ✅"""
        key_bytes = str(user_id).encode('utf-8')
        hash_value = int(hashlib.md5(key_bytes).hexdigest(), 16)
        return hash_value % self.num_shards   # Returns 0, 1, 2, or 3

    def execute(self, user_id: int, query: str, params: tuple = ()):
        """Run SQL on the correct shard — app never needs to know which one."""
        shard_idx = self.get_shard(user_id)
        config = self.shards[shard_idx]
        print(f"→ user_id={user_id} → Shard {shard_idx} ({config['host']})")

        conn = psycopg2.connect(**config, user="app_user", password="secret")
        try:
            cur = conn.cursor()
            cur.execute(query, params)
            conn.commit()
            return cur.fetchall()
        finally:
            conn.close()

# ── Usage ────────────────────────────────────────────────────────────
router = ShardRouter()

# Write: insert an order for user 8,500,000
router.execute(
    user_id=8_500_000,
    query="INSERT INTO orders(user_id, item) VALUES (%s, %s)",
    params=(8_500_000, "Margherita Pizza")
)
# Output: → user_id=8500000 → Shard 2 (shard-c.orderease.com)

# Read: get all orders for user 1,234,567
results = router.execute(
    user_id=1_234_567,
    query="SELECT * FROM orders WHERE user_id = %s",
    params=(1_234_567,)
)
# Output: → user_id=1234567 → Shard 1 (shard-b.orderease.com)

⚠️ The Hotspot Problem — And How to Fix It

A hotspot happens when one shard gets far more traffic than the others.
It overloads → queries slow → server crashes. 🔥
This is the #1 real-world sharding failure.

🔴 Hotspot — India Launch Day (20M signups in 1 hour)

Shard A 🔴
ALL new users!
87% load 💥
Shard B ✅
idle
5% load
Shard C ✅
idle
5% load
Shard D ✅
idle
3% load

Root cause: range sharding put all new Indian user_ids into Shard A's range.
Fix: switch to hash-based sharding → even spread across all 4 shards ✅


🔑 Choosing the Right Shard Key

┌─────────────────┬───────────────────────────────────────────────┐
│  Candidate Key  │  Verdict                                      │
├─────────────────┼───────────────────────────────────────────────┤
│  created_at     │ ❌ TERRIBLE — all new writes hit one shard    │
│  restaurant_id  │ ❌ BAD — popular restaurants = hotspot         │
│  city           │ ⚠️  RISKY — Mumbai gets 10x more traffic      │
│  user_id (hash) │ ✅ GREAT — even, high cardinality 🏆           │
│  order_id (UUID)│ ✅ GREAT — random, balanced perfectly          │
└─────────────────┴───────────────────────────────────────────────┘

  Golden Rule: High cardinality + used in most queries + NOT sequential

🌍 Full OrderEase Production Architecture

           ORDEREASE APP (500M users)
                     │
        ┌────────────▼────────────┐
        │   Redis Cache Layer     │ ← Absorbs 85% of reads
        └────────────┬────────────┘
                     │ cache miss
        ┌────────────▼────────────┐
        │  Consistent Hash Router │ ← ProxySQL / Vitess
        └──┬──────┬──────┬──────┬─┘
           │      │      │      │
    ┌──────▼┐ ┌───▼──┐ ┌─▼────┐ ┌▼──────┐
    │Shard A│ │Shard B│ │Shard C│ │Shard D│
    │+Repli.│ │+Repli.│ │+Repli.│ │+Repli.│ ← Active Data Guard
    └───────┘ └───────┘ └───────┘ └───────┘
                     │
        ┌────────────▼────────────┐
        │  Analytics DB (BigQuery)│ ← Cross-shard queries, dashboards
        │  Fed via GoldenGate     │   ML model training
        └─────────────────────────┘
✅ When should you actually shard?

🟢 Single DB above 500GB and queries are visibly slow.
🟢 Write throughput exceeds 50,000 writes/second on one machine.
🟢 Data residency laws require storing data in specific countries.

🔴 Try these FIRST — they buy months of headroom:
→ Better indexes (fixes 80% of slowness cheaply)
→ Read replicas (offload reads immediately)
→ Redis caching (eliminates most DB reads entirely)
→ Vertical scaling (bigger server — far simpler than sharding)
🚫 Never Do These With Sharding:

❌ Never use a timestamp as your shard key — instant hotspot on every new day.
❌ Never use AUTO_INCREMENT IDs — each shard produces id=1, id=2 — collisions everywhere. Use UUIDs.
❌ Never add a UNIQUE constraint across shards — uniqueness only enforces within one shard.
❌ Never shard prematurely — sharding adds enormous operational complexity.

Happy sharding! 🔀✨

Comments