β‘ MongoDB Sharding Cheatsheet
Complete guide to horizontal scaling, shard keys, zones, and production best practices
What is Sharding?
Horizontal scaling across multiple machines
Sharding is MongoDB's approach to horizontal scaling - splitting your data across multiple servers (shards) to handle massive datasets and high throughput.
Think of it like organizing a massive library: instead of one giant building, you have multiple branches, each holding different books based on categories.
β’ Limited by hardware
β’ Downtime for upgrades
β’ Exponentially expensive
β’ Single point of failure
β’ Near-infinite scalability
β’ No downtime
β’ Cost-effective (commodity hardware)
β’ High availability built-in
Sharding Architecture
Components of a sharded cluster
| Component | Role | Count | Notes |
|---|---|---|---|
| mongos | Query router - directs operations to shards | 1+ (multiple for HA) | Stateless, lightweight |
| Config Servers | Store cluster metadata & chunk mapping | 3 (replica set) | Critical - always 3! |
| Shards | Store actual data (replica sets) | 2+ recommended | Each is a replica set |
| Chunks | Data ranges (64MB by default) | Auto-managed | Balanced automatically |
Shard Key Strategies
Three main approaches
| Strategy | Distribution | Range Queries | Best For |
|---|---|---|---|
| Hashed | β Perfect | β No | Monotonic IDs, random access |
| Ranged | β οΈ Varies | β Yes | Time-series, geographic, categories |
| Compound | β Good | β On prefix | Multi-tenant, balanced workloads |
Zone Sharding
Geographic or logical data placement
Chunk Balancing
Automatic data distribution
MongoDB automatically balances chunks across shards. The balancer runs in the background and moves chunks when imbalance is detected (typically when difference > 8 chunks).
- Schedule balancing during off-peak hours
- Stop balancer before maintenance or backups
- Monitor chunk distribution regularly
- Jumbo chunks (>64MB) won't be moved - investigate!
Monitoring & Diagnostics
Essential commands
Production Best Practices
Critical guidelines
- Test with production workload - Simulate queries before sharding
- Check cardinality - Unique values > number of shards Γ 10
- Avoid monotonic keys - timestamp, _id, auto-increment
- Consider query patterns - What fields are in WHERE clause?
- Plan for growth - Will key work at 10x, 100x scale?
- Use compound keys - For better targeting
- Pre-split chunks - For known distribution
- Monitor hotspots - Use sh.status() regularly
- Using timestamp as shard key - All writes go to one shard!
- Low cardinality keys - gender, status, type (few unique values)
- Not testing query patterns - Results in scatter-gather queries
- Forgetting about zones - Missing compliance/performance opportunities
- Not monitoring balancer - Uneven distribution goes unnoticed
| Scenario | Recommended Shard Key | Why |
|---|---|---|
| User documents | { _id: "hashed" } |
Even distribution, ObjectId is monotonic |
| Orders by customer | { customerId: 1 } |
Query by customer, high cardinality |
| Time-series logs | { serverId: 1, timestamp: 1 } |
Avoid timestamp-only hotspot |
| Multi-tenant SaaS | { tenantId: 1, _id: 1 } |
Tenant isolation + even distribution |
| E-commerce products | { category: 1, productId: 1 } |
Category queries + uniqueness |
| Social media posts | { userId: "hashed" } |
User-centric, even distribution |
Don't shard too early! Sharding adds complexity. Consider sharding when:
- Dataset > 2-3 TB per server
- Working set > RAM (performance degradation)
- Write throughput > single server capacity
- Need geographic data distribution
Before sharding: Try vertical scaling, indexing optimization, replica sets with read preference