Section 14: Cheatsheets

⚑ MongoDB Sharding Cheatsheet

Complete guide to horizontal scaling, shard keys, zones, and production best practices

πŸ“š

What is Sharding?

Horizontal scaling across multiple machines

πŸ’‘ Sharding in Simple Terms

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.

πŸ—οΈ Sharded vs Non-Sharded Architecture
❌ Without Sharding Single Server 100 GB Data 10K ops/sec ❌ Single Point of Failure Problems: β€’ Limited by single server β€’ Can't scale beyond hardware β€’ Downtime = full outage βœ… With Sharding Shard 1 33 GB A-F data Shard 2 33 GB G-M data Shard 3 34 GB N-Z data mongos (Router) Benefits: β€’ 3x capacity (300 GB total) β€’ 3x throughput (30K ops/sec) β€’ High availability β€’ Horizontal scaling
Vertical Scaling (Expensive!)
β€’ Buy bigger server
β€’ Limited by hardware
β€’ Downtime for upgrades
β€’ Exponentially expensive
β€’ Single point of failure
Horizontal Scaling (Sharding!)
β€’ Add more servers
β€’ Near-infinite scalability
β€’ No downtime
β€’ Cost-effective (commodity hardware)
β€’ High availability built-in
πŸ—οΈ

Sharding Architecture

Components of a sharded cluster

πŸ”§ Sharded Cluster Components
Application mongos mongos mongos Config Servers (Metadata) β€’ Shard map β€’ Chunks Shard 1 (Replica Set) Primary Secondary Secondary Chunks: 0-33 Shard 2 (Replica Set) Primary Secondary Secondary Chunks: 34-66
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 Keys

The most critical decision in sharding

⚠️ CRITICAL: Shard Key Cannot Be Changed! Once you shard a collection, the shard key is immutable! Choose wisely or face a complete data migration.
Good Shard Key Qualities
Essential
A good shard key has three critical properties:
βœ… The 3 Golden Rules
  1. High Cardinality - Many unique values (userId βœ“, gender βœ—)
  2. Low Frequency - Values evenly distributed (not 90% same value)
  3. Non-Monotonic - Avoid always-increasing values (timestamp βœ—)
Enable Sharding
Setup
mongosh (connect to mongos)
> // Step 1: Enable sharding on database > sh.enableSharding("ecommerce")
{ ok: 1, '$clusterTime': { ... } }
> // Step 2: Shard the collection > sh.shardCollection( "ecommerce.orders", { customerId: 1 } // Shard key )
{ collectionsharded: 'ecommerce.orders', ok: 1 }
> // Verify sharding status > sh.status()
🎯

Shard Key Strategies

Three main approaches

Hashed Shard Key
Best for Monotonic
Hash the field value to distribute data evenly. Perfect for monotonically increasing values.
Example: User IDs (ObjectId)
> // Create hashed index first > db.users.createIndex({ _id: "hashed" }) > // Shard on hashed _id > sh.shardCollection( "myapp.users", { _id: "hashed" } // Hashed shard key! ) > // Result: Even distribution across shards > // _id: ObjectId("507f...") β†’ hash β†’ Shard 2 > // _id: ObjectId("508a...") β†’ hash β†’ Shard 1 > // _id: ObjectId("509b...") β†’ hash β†’ Shard 3
Advantages
β€’ Even distribution guaranteed
β€’ No hotspots
β€’ Works with monotonic values
β€’ Simple to implement
Disadvantages
β€’ No range queries!
β€’ Can't use find({ _id: { $gt: ... }})
β€’ Broadcasts scatter-gather queries
β€’ Less predictable data location
Ranged Shard Key
Best for Queries
Use natural field values for sharding. Great for range queries, but watch for hotspots!
Example: Geographic Region
> // Shard by country > sh.shardCollection( "ecommerce.customers", { country: 1 } // Ranged shard key ) > // Result: Geographic clustering > // Shard 1: USA customers > // Shard 2: UK, France, Germany > // Shard 3: India, China, Japan > // βœ… Efficient range query (targets one shard!) > db.customers.find({ country: "USA" }) > // ⚠️ Watch out for hotspots! > // If 80% users are from USA β†’ Shard 1 overloaded
Advantages
β€’ Efficient range queries
β€’ Targeted reads/writes
β€’ Predictable data location
β€’ Good for time-series
Risks
β€’ Possible hotspots
β€’ Uneven distribution
β€’ Monotonic values = one shard gets all writes
β€’ Needs careful key selection
Compound Shard Key
Best Balance
Combine multiple fields for better distribution and query targeting.
Example: Region + User ID
> // Compound shard key: region + userId > sh.shardCollection( "social.posts", { region: 1, userId: 1 } // Compound! ) > // Benefits: > // 1. Region provides coarse distribution > // 2. userId provides fine-grained distribution > // βœ… Targeted query (one shard) > db.posts.find({ region: "us-east", userId: 12345 }) > // βœ… Reasonable query (subset of shards) > db.posts.find({ region: "us-east" }) > // ❌ Scatter-gather (all shards) > db.posts.find({ userId: 12345 }) // Missing prefix!
πŸ’‘ Compound Shard Key Rules
  • First field should have low cardinality (10-100 values) for distribution
  • Second field should have high cardinality for uniqueness
  • Queries must include the first field to target specific shards
  • Example combos: {region: 1, customerId: 1}, {category: 1, productId: 1}
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

Zone-Based Sharding
Data Locality
Assign data ranges to specific shards based on geography or compliance requirements.
Example: GDPR Compliance
> // Define zones for shards > sh.addShardToZone("shard01", "EU") > sh.addShardToZone("shard02", "US") > sh.addShardToZone("shard03", "ASIA") > // Associate data ranges with zones > sh.updateZoneKeyRange( "users.customers", { country: "Austria" }, // Min { country: "UK" }, // Max (exclusive) "EU" // Zone ) > sh.updateZoneKeyRange( "users.customers", { country: "USA" }, { country: "USA" }, "US" ) > // Result: EU data stays in EU shards (GDPR compliant!) > // US data stays in US shards
βœ… Zone Sharding Use Cases
  • Compliance: Keep EU data in EU (GDPR), China data in China
  • Performance: Data close to users (low latency)
  • Hot/Cold data: Recent data on SSD shards, old data on HDD
  • Multi-tenant: Premium customers on high-performance shards
βš–οΈ

Chunk Balancing

Automatic data distribution

πŸ’‘ How Balancing Works

MongoDB automatically balances chunks across shards. The balancer runs in the background and moves chunks when imbalance is detected (typically when difference > 8 chunks).

Balancer Management
> // Check balancer status > sh.getBalancerState()
true
> // Check if balancer is running > sh.isBalancerRunning()
false
> // Stop balancer (for maintenance) > sh.stopBalancer() > // Start balancer > sh.startBalancer() > // Set balancer window (off-peak hours) > db.settings.updateOne( { _id: "balancer" }, { $set: { activeWindow: { start: "23:00", // 11 PM stop: "06:00" // 6 AM } } }, { upsert: true } )
⚠️ Balancer Best Practices
  • 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

Sharding Status Commands
> // Complete cluster status > sh.status() > // Check collection sharding details > db.collection.getShardDistribution()
Shard shard01 at shard01/... data: 512 MiB docs: 1024000 chunks: 42 estimated data per chunk: 12 MiB estimated docs per chunk: 24380 Shard shard02 at shard02/... data: 498 MiB docs: 996000 chunks: 40 estimated data per chunk: 12 MiB estimated docs per chunk: 24900 Totals data: 1010 MiB docs: 2020000 chunks: 82
> // Check for jumbo chunks > db.chunks.find({ jumbo: true }) > // View chunk ranges > db.chunks.find({ ns: "mydb.mycollection" }) .sort({ min: 1 }) .pretty()
⭐

Production Best Practices

Critical guidelines

βœ… Shard Key Selection Checklist
  1. Test with production workload - Simulate queries before sharding
  2. Check cardinality - Unique values > number of shards Γ— 10
  3. Avoid monotonic keys - timestamp, _id, auto-increment
  4. Consider query patterns - What fields are in WHERE clause?
  5. Plan for growth - Will key work at 10x, 100x scale?
  6. Use compound keys - For better targeting
  7. Pre-split chunks - For known distribution
  8. Monitor hotspots - Use sh.status() regularly
❌ Common Sharding Mistakes
  • 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
πŸ’‘ When to Shard?

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