Algocraft
/
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
↓ ↓ ↓
🔥 Hot Shard Detected!
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 —
3×
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