/ System Design / Database Sharding ← All Designs
🔀

Database Sharding

Horizontally partition data across multiple database nodes to scale beyond one machine

Horizontal Scaling Partitioning MongoDB / Cassandra Hot Shard Problem
Shard Distribution Hash Sharding
SHARD ROUTER
hash(key) % N
↓ ↓ ↓
Record Distribution (Bar Chart)
⚠️
Cross-Shard Query Cost
Queries that don't include the shard key must fan out to ALL 3 shards, then merge results — the work of a single-shard query.

Key Concepts

🎯 When to Use
  • Single DB can't handle data volume
  • Write throughput exceeds one node
  • Data too large for vertical scaling
  • Geographic data partitioning needed
⚖️ Trade-offs
  • Cross-shard queries are expensive
  • Adding shards requires resharding
  • Hot shards cause uneven load
  • Transactions span multiple shards
🎓 Interview Tips
  • Choose shard key carefully — no changes later
  • Consistent hashing avoids full resharding
  • Vitess / CockroachDB handle resharding
  • Avoid shard keys that create hot spots