Database Sharding in System Design: Strategies, Shard Keys and Scaling Explained
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. 🔀
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.
🗂️ 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
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
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
~4ms latency ✅
~6ms latency ✅
~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
"Get order for user_id = 8,500,000"
hash(8500000) % 4 = 2 → Shard C
Searches 125M rows — not 500M! → returns in 4ms ✅
10× faster than a single un-sharded database
💻 Building a Shard Router in Python
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)
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
└─────────────────────────┘
🟢 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 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
Post a Comment