Section 7: Advanced MongoDB

⚑ MongoDB Sharding

Master horizontal scaling by distributing data across multiple servers for infinite scalability

πŸ“– The Viral App Crisis: Mike's Sharding Journey

😱 Day 1: The Meltdown

Mike built a social media app called "ShareSnap" that went VIRAL overnight. Featured on TechCrunch, users flooded in. From 10 million to 100 million users in just 3 weeks! His single MongoDB server was completely overwhelmed:

πŸ’Ύ
800GB
Data size, growing 50GB daily!
⏱️
30s
Query response time
πŸ”₯
95%
RAM usage (swapping to disk)

Users were complaining: "App is so slow!", "I can't upload photos!", "Posts take forever to load!" Mike was panicking. His investors were worried. His server was literally melting under the load! πŸ”₯

πŸ€” Day 2: Evaluating Options

Mike's senior developer suggested two options:

❌ Option 1: Vertical Scaling

"Buy ONE huge server"

  • Specs: 128GB RAM, 2TB SSD, 32 CPU cores
  • Cost: $8,000/month
  • Problem: Still has limits! What if we grow to 500M users?
  • Risk: Single point of failure
  • Scalability: Limited to hardware constraints

βœ… Option 2: Horizontal Scaling (Sharding)

"Split data across MULTIPLE servers"

  • Specs: 4 servers Γ— 32GB RAM, 500GB each
  • Cost: $6,000/month total
  • Bonus: Can add MORE servers anytime!
  • Reliability: No single point of failure
  • Scalability: Infinite! Just add servers

πŸ’‘ Mike's Decision: "Let's go with sharding! It's cheaper NOW and will scale infinitely. Plus, if one server fails, the others keep running!"

πŸŽ‰ Day 3: After Implementing Sharding

Mike's team sharded the database by userId across 4 servers. Each shard handled 25 million users and 200GB of data.

⚑
300ms
Query time (100x faster!)
πŸ“ˆ
∞
Infinitely scalable
πŸ’°
$2K
Monthly savings
😊
4.8β˜…
App store rating

πŸš€ 6 Months Later: 500 Million Users!

ShareSnap grew to 500 million users. Mike simply added 12 more shards (16 total). No downtime. No migration nightmare. Just smooth horizontal scaling. Total cost: $24,000/month for 500M users vs $50,000+ for trying to vertically scale!

πŸ’‘ Key Lesson: Vertical scaling (bigger servers) has limits. Horizontal scaling (sharding) lets you grow infinitely by adding more machines!

πŸ”€ What is Sharding?

Sharding is MongoDB's horizontal scaling method where data is distributed across multiple servers (called shards). Each shard stores a portion of your data. Together, they act as one logical database, allowing your application to handle massive datasets and traffic.

πŸ“Š Single Server vs Sharded Cluster

❌ Single Server (Vertical Scaling)

ONE HUGE SERVER
128GB RAM | 2TB Storage

Contains:

πŸ”΄ ALL 100M users
πŸ”΄ ALL 800GB data
πŸ”΄ ALL queries hit here

Problems:

  • Slow (30s queries)
  • RAM exhausted
  • Can't scale infinitely
  • Single point of failure
  • Very expensive ($8K/mo)

βœ… Sharded Cluster (Horizontal Scaling)

Shard 1
8GB RAM
Users 0-25M
200GB
Shard 2
8GB RAM
Users 25M-50M
200GB
Shard 3
8GB RAM
Users 50M-75M
200GB
Shard 4
8GB RAM
Users 75M-100M
200GB

Benefits:

  • ⚑ Fast (300ms queries)
  • πŸ“ˆ Infinitely scalable
  • πŸ’° Cheaper ($6K/mo)
  • πŸ›‘οΈ No single point of failure
  • 🌍 Geographic distribution

πŸ”‘ Key Insight: In sharding, each server (shard) only handles a portion of the data. This means queries are faster (less data to scan), RAM usage is lower (smaller working set), and you can add more servers whenever needed!

🎯 Simple Analogy: Library Books

πŸ“š Without Sharding:

Imagine ONE tiny library room with 1 million books stacked floor to ceiling. Finding a book takes FOREVER. The librarian is overwhelmed. Adding more books? Impossible!

πŸ“– With Sharding:

Now imagine 10 library rooms. Books A-C in Room 1, D-F in Room 2, etc. Finding a book is FAST (you know which room to go to). Each librarian manages fewer books. Need more capacity? Add more rooms!

πŸ€” Why Do We Need Sharding?

Sharding solves three critical problems that occur when your application grows:

πŸ’Ύ Problem 1: Data Storage Limit

Scenario: Your e-commerce app has 100 million products. Each product has images, reviews, pricing history. Your database is now 5TB!

Problem: A single server can't hold 5TB in memory (RAM). It has to constantly read from disk, making everything SLOW.

βœ… Solution with Sharding: Split 5TB across 10 servers = 500GB per server. Each shard only needs 64GB RAM to keep its working set in memory!

πŸš€ Problem 2: Throughput Limit

Scenario: Your social media app has 50 million active users. During peak hours, you get 100,000 queries per second!

Problem: A single server can only handle ~10,000 queries/sec (depending on complexity). Requests queue up. Users see "Loading..." forever.

βœ… Solution with Sharding: 10 shards Γ— 10,000 queries/sec = 100,000 queries/sec total capacity! Each shard handles its portion of users.

🌍 Problem 3: Geographic Distribution

Scenario: You have users in India, USA, and Europe. They all connect to your server in Mumbai. US/Europe users experience 300ms+ latency.

Problem: Physics! Data takes time to travel 12,000 km. No amount of server upgrades can reduce network latency.

βœ… Solution with Sharding: Place shards geographically. Indian users β†’ Mumbai shard, US users β†’ Virginia shard, EU users β†’ Frankfurt shard. Latency drops to <50ms!

πŸ“Š When Should You Consider Sharding?

πŸ“ˆ

Data Size: Database > 200GB and growing

⚑

Traffic: > 10,000 operations/sec

πŸ’°

Cost: Vertical scaling too expensive

🌍

Geography: Global user base

πŸ“– The Viral App Crisis

😱 Mike's Database Meltdown

Mike's social media app went viral overnight. Users flooded in: 10 million to 100 million in 3 weeks. His single MongoDB server (32GB RAM, 500GB storage) was dying. Queries taking 30+ seconds. Database size: 800GB and growing 50GB/day!

πŸ’Ύ
800GB
Database Size
Growing 50GB daily
⏱️
30s
Query Time
Users leaving!
πŸ”₯
95%
RAM Usage
Constant swapping

Mike's Options:
1️⃣ Vertical Scaling: Bigger server (128GB RAM, 2TB storage) = $8,000/month
2️⃣ Horizontal Scaling (Sharding): Split across 4 servers = $6,000/month total

😊 After Implementing Sharding

Mike implemented sharding across 4 servers. Data automatically distributed by user region. Each shard handled 25 million users and 200GB data.

⚑
300ms
Query Time
100x faster! πŸš€
πŸ“ˆ
∞
Scalability
Add shards anytime
πŸ’°
$2K
Monthly Savings
vs vertical scaling

πŸ’‘ Key Lesson: Vertical scaling (bigger servers) has limits. Horizontal scaling (sharding) lets you grow infinitely by adding more machines!

πŸ”€ What is Sharding?

Sharding is a method of distributing data across multiple machines. MongoDB splits your data into smaller chunks and distributes them across multiple servers called shards. Each shard is an independent database.

πŸ“Š Non-Sharded vs Sharded Architecture

❌ Single Server (Non-Sharded) Your Application MongoDB Server 32GB RAM | 500GB Storage ALL 100M Users β€’ User 1 - 25M β€’ User 25M - 50M β€’ User 50M - 75M β€’ User 75M - 100M 800GB Total Data ⚠️ RAM exhausted ⚠️ Slow queries (30s) ⚠️ Can't scale further VS βœ… Sharded Cluster Your Application mongos (Query Router) Shard 1 8GB RAM | 250GB User 1-25M Region: US-East 200GB data Shard 2 8GB RAM | 250GB User 25M-50M Region: US-West 200GB data + Shard 3 (Europe) + Shard 4 (Asia) βœ… 4 Servers Total βœ“ Fast queries (300ms) βœ“ Infinitely scalable βœ“ Cost effective

🎯 Why Shard?

πŸ“ˆ

Horizontal Scaling

Add more machines to handle growing data. Much cheaper than buying bigger servers!

⚑

Better Performance

Queries distributed across shards run in parallel. 4 shards = 4x throughput!

πŸ’Ύ

More Storage

Each shard has its own storage. 4 shards with 1TB each = 4TB total capacity!

🌍

Geographic Distribution

Place shards closer to users. US users β†’ US shard, Asia users β†’ Asia shard!

πŸ”„

High Availability

If one shard fails, others continue working. Your app stays online!

πŸ’°

Cost Effective

4 medium servers cost less than 1 massive server. Scale economically!

πŸ“Š When to Use Sharding?

βœ… Good Candidates:
  • Working set exceeds RAM (>64GB)
  • Database size > 1TB
  • Write-heavy workload
  • Geographic distribution needed
  • Predictable growth pattern
❌ Not Good For:
  • Small databases (<100GB)
  • Low traffic applications
  • Complex join-heavy queries
  • Unpredictable access patterns
  • Budget constraints (need 3+ servers)

πŸ—οΈ Sharding Architecture

A MongoDB sharded cluster consists of three main components working together to distribute and manage your data across multiple servers.

πŸ“ Complete Sharded Cluster Architecture

Application Servers App 1 App 2 App 3 Query Routers (mongos) mongos 1 Routes queries mongos 2 Routes queries mongos 3 Routes queries Config Servers (Metadata) Config Server Replica Set Stores: Chunk ranges, Shard locations, Cluster metadata, Zone configurations Shards (Data Storage) Shard 1 (Replica Set) Primary Read/Write Secondary Read only Sec Backup Data Chunks Chunk 1: userId 0-25M Chunk 3: userId 50M-75M 200GB data Shard 2 (Replica Set) Primary Read/Write Secondary Read only Sec Backup Data Chunks Chunk 2: userId 25M-50M Chunk 4: userId 75M-100M 200GB data + More Shards... Scale infinitely!

πŸ”§ Core Components

πŸ”€

mongos (Query Router)

The mongos acts as the interface between your application and the sharded cluster. It routes queries to the appropriate shards and merges results.

Key Functions:
  • Query routing: Determines which shard(s) to query
  • Result merging: Combines results from multiple shards
  • Write distribution: Routes inserts/updates to correct shards
  • Metadata caching: Caches chunk locations for faster routing
βš™οΈ

Config Servers (Metadata Store)

Config servers store all the metadata about the sharded cluster - which chunks are on which shards, shard locations, and cluster configuration.

Stored Metadata:
  • Chunk ranges: Which data ranges (chunks) exist
  • Chunk locations: Which shard contains each chunk
  • Shard information: List of all shards in the cluster
  • Database/collection: Which collections are sharded
  • Zone mappings: Geographic or custom zone configurations
⚠️ Critical: Always deployed as a replica set (3+ servers) for high availability. If config servers go down, the cluster becomes read-only!
πŸ’Ύ

Shards (Data Storage)

Each shard is a MongoDB replica set that stores a subset of the sharded data. Data is distributed across shards based on the shard key.

Shard Characteristics:
  • Replica set: Each shard is a full replica set (primary + secondaries)
  • Data subset: Contains only part of the total data (specific chunks)
  • Independent: Can operate independently if other shards fail
  • Scalable: Add more shards as data grows
  • Balanced: MongoDB automatically balances chunks across shards

πŸ”„ How They Work Together

  1. Application sends query to mongos
  2. mongos checks config servers for chunk locations
  3. mongos routes query to appropriate shard(s)
  4. Shard(s) execute query and return results
  5. mongos merges results and returns to application

πŸ”‘ Shard Keys: The Most Critical Decision

⚠️ WARNING: The shard key is IMMUTABLE after choosing it. Choose wisely - you can't change it without rebuilding the entire collection!

A shard key is the indexed field(s) that MongoDB uses to partition data across shards. It determines which shard stores which documents. Think of it as the "address" that tells MongoDB where to store and find each document.

πŸ“Š How Shard Key Distributes Data

Documents with userId (Shard Key) { userId: 1000 } name: "Alice" email: "alice@..." { userId: 50000000 } name: "Bob" email: "bob@..." { userId: 95000000 } name: "Charlie" email: "charlie@..." MongoDB Evaluates userId Distributed Across Shards Shard 1 Range: userId 0 - 33,000,000 βœ“ Alice (userId: 1000) Shard 2 Range: userId 33M - 66M βœ“ Bob (userId: 50M) Shard 3 Range: userId 66M - 100M βœ“ Charlie (userId: 95M)

βœ… Good vs ❌ Bad Shard Keys

βœ… Good Shard Keys

  • High cardinality: Many unique values (userId, email, orderId)
  • Even distribution: Values spread evenly (hashed userId)
  • Query isolation: Queries target single shard (region + userId)
  • Monotonic writes avoided: Not always increasing (_id, timestamp ❌)
  • Frequently queried: Used in most queries
Examples:
{ userId: 1 }
{ email: "hashed" }
{ region: 1, userId: 1 }

❌ Bad Shard Keys

  • Low cardinality: Few unique values (gender, status, boolean)
  • Uneven distribution: Most docs in one range (VIP status)
  • Monotonic: Always increasing (timestamp, _id, date)
  • Not in queries: Never used in WHERE clauses
  • Single hot shard: All writes to one shard
Examples:
{ _id: 1 } ❌
{ createdAt: 1 } ❌
{ status: 1 } ❌

🎯 Shard Key Selection Criteria

πŸ“Š

High Cardinality

Many unique values ensure data spreads across shards. Million users = great cardinality!

βš–οΈ

Even Distribution

Values distributed evenly prevent hot shards. Hashed keys guarantee this!

🎯

Query Isolation

Queries should target one shard, not broadcast to all shards. Faster queries!

πŸ”

Used in Queries

Shard key should be in most WHERE clauses to enable targeted queries!

Creating a Sharded Collection
// Step 1: Enable sharding on database
sh.enableSharding("myDatabase")

// Step 2: Create index on shard key (REQUIRED before sharding)
db.users.createIndex({ userId: 1 })

// Step 3: Shard the collection
sh.shardCollection("myDatabase.users", { userId: 1 })

// For compound shard key:
db.users.createIndex({ region: 1, userId: 1 })
sh.shardCollection("myDatabase.users", { region: 1, userId: 1 })

// For hashed shard key (better distribution):
db.users.createIndex({ userId: "hashed" })
sh.shardCollection("myDatabase.users", { userId: "hashed" })

// Verify sharding status
sh.status()

πŸ“¦ Chunks: How Data is Split

MongoDB doesn't distribute individual documentsβ€”it distributes chunks. A chunk is a contiguous range of shard key values. Default chunk size is 128MB (configurable).

πŸ“Š Chunk Distribution Across Shards

Shard Key Range: userId 0 to 100,000,000 0 100M 50M Divided into Chunks (128MB each) Chunk 1 Range: 0 - 25M Size: 128MB Chunk 2 Range: 25M - 50M Size: 128MB Chunk 3 Range: 50M - 75M Size: 128MB Chunk 4 Range: 75M - 100M Size: 128MB Distributed Across Shards Shard 1 Chunk 1 0 - 25M 128MB + More chunks... Shard 2 Chunk 2 25M - 50M 128MB + More chunks... Shard 3 Chunk 3 & 4 50M - 100M 256MB + More chunks...

πŸ”„ Chunk Splitting & Migration

βœ‚οΈ Chunk Splitting

When a chunk exceeds 128MB, MongoDB automatically splits it into two smaller chunks.

Example:
Chunk: 0-50M (150MB)
β†’ Split into:
β€’ Chunk A: 0-25M (75MB)
β€’ Chunk B: 25M-50M (75MB)

πŸ”„ Chunk Migration

The balancer automatically moves chunks between shards to keep data evenly distributed.

Example:
Shard 1: 100 chunks
Shard 2: 50 chunks
β†’ Balancer moves 25 chunks
β†’ Result: 75 chunks each βœ“

πŸ’‘ Key Chunk Facts

  • Default size: 128MB (configurable: 1MB - 1024MB)
  • Splitting: Automatic when chunk exceeds size limit
  • Migration: Automatic via balancer (runs in background)
  • Indivisible: A chunk lives entirely on one shard
  • Metadata: Chunk ranges stored in config servers
  • Performance: Smaller chunks = more frequent migrations but better distribution

🎯 Sharding Strategies

MongoDB offers three main sharding strategies. Each has different data distribution patterns and use cases.

πŸ“Š 1. Ranged Sharding

Divides data based on ranges of shard key values. Documents with "nearby" shard key values are likely on the same shard.

Ranged Sharding: userId 0-100M Shard 1 Range: 0 - 33M Documents: β€’ userId: 100 β€’ userId: 5000 β€’ userId: 25M Shard 2 Range: 33M - 66M Documents: β€’ userId: 35M β€’ userId: 50M β€’ userId: 65M Shard 3 Range: 66M - 100M Documents: β€’ userId: 70M β€’ userId: 85M β€’ userId: 99M
βœ… Pros:
  • Range queries very efficient
  • Good for sorted data access
  • Related data stays together
❌ Cons:
  • Risk of hot spots
  • Uneven distribution possible
  • Monotonic keys = single shard writes

#️⃣ 2. Hashed Sharding

Uses a hash function on shard key values to distribute data. Guarantees even distribution but loses range query efficiency.

Hashed Sharding: hash(userId) Hash Function userId: 100 β†’ hash: 7384... userId: 50M β†’ hash: 2891... Shard 1 β€’ userId: 50M (hashed) β€’ userId: 200 (hashed) β€’ userId: 99M (hashed) ~33% of data Shard 2 β€’ userId: 100 (hashed) β€’ userId: 75M (hashed) β€’ userId: 1000 (hashed) ~33% of data Shard 3 β€’ userId: 25M (hashed) β€’ userId: 5000 (hashed) β€’ userId: 80M (hashed) ~33% of data
βœ… Pros:
  • Perfect even distribution
  • No hot spots
  • Works with monotonic keys
❌ Cons:
  • Range queries = broadcast
  • Related data scattered
  • Can't use for sorting

🌍 3. Zone/Tag-Aware Sharding

Assigns data to specific shards based on zones. Perfect for geographic distribution or compliance requirements.

Zone Sharding by Region πŸ‡ΊπŸ‡Έ US Shard Zone: "US" region: "US" β€’ New York users β€’ LA users β€’ Chicago users Server in Virginia πŸ‡ͺπŸ‡Ί EU Shard Zone: "EU" region: "EU" β€’ London users β€’ Berlin users β€’ Paris users Server in Frankfurt 🌏 Asia Shard Zone: "ASIA" region: "ASIA" β€’ Tokyo users β€’ Singapore users β€’ Mumbai users Server in Singapore
βœ… Pros:
  • Data locality (low latency)
  • Compliance (GDPR, data residency)
  • Custom distribution rules
❌ Cons:
  • More complex setup
  • Manual zone management
  • Uneven loads across zones

βš™οΈ Setting Up a Sharded Cluster

Here's a complete step-by-step guide to set up MongoDB sharding from scratch.

Step 1: Start Config Server Replica Set
# Start 3 config servers (always use replica set)
mongod --configsvr --replSet configRS --port 27019 --dbpath /data/configdb1
mongod --configsvr --replSet configRS --port 27020 --dbpath /data/configdb2
mongod --configsvr --replSet configRS --port 27021 --dbpath /data/configdb3

# Connect and initiate replica set
mongosh --port 27019
rs.initiate({
  _id: "configRS",
  configsvr: true,
  members: [
    { _id: 0, host: "localhost:27019" },
    { _id: 1, host: "localhost:27020" },
    { _id: 2, host: "localhost:27021" }
  ]
})
Step 2: Start Shard Replica Sets
# Start Shard 1 replica set
mongod --shardsvr --replSet shard1RS --port 27022 --dbpath /data/shard1a
mongod --shardsvr --replSet shard1RS --port 27023 --dbpath /data/shard1b
mongod --shardsvr --replSet shard1RS --port 27024 --dbpath /data/shard1c

# Initiate Shard 1
mongosh --port 27022
rs.initiate({
  _id: "shard1RS",
  members: [
    { _id: 0, host: "localhost:27022" },
    { _id: 1, host: "localhost:27023" },
    { _id: 2, host: "localhost:27024" }
  ]
})

# Repeat for Shard 2, Shard 3, etc...
Step 3: Start mongos Router
# Start mongos (query router)
mongos --configdb configRS/localhost:27019,localhost:27020,localhost:27021 --port 27017

# mongos is now running on port 27017
# Your application connects to mongos, not directly to shards!
Step 4: Add Shards to Cluster
# Connect to mongos
mongosh --port 27017

# Add each shard
sh.addShard("shard1RS/localhost:27022,localhost:27023,localhost:27024")
sh.addShard("shard2RS/localhost:27025,localhost:27026,localhost:27027")
sh.addShard("shard3RS/localhost:27028,localhost:27029,localhost:27030")

# Verify shards added
sh.status()
Step 5: Enable Sharding & Shard Collection
# Enable sharding on database
sh.enableSharding("myDatabase")

# Create index on shard key (REQUIRED!)
db.users.createIndex({ userId: 1 })

# Shard the collection with ranged sharding
sh.shardCollection("myDatabase.users", { userId: 1 })

# OR use hashed sharding for better distribution
db.users.createIndex({ userId: "hashed" })
sh.shardCollection("myDatabase.users", { userId: "hashed" })

# Verify sharding
db.users.getShardDistribution()

πŸŽ‰ Congratulations! Your sharded cluster is now running. MongoDB will automatically distribute data across shards and balance chunks.

🌍 Setting Up Zone Sharding

Configure Geographic Zones
# Tag shards with zones
sh.addShardTag("shard1RS", "US")
sh.addShardTag("shard2RS", "EU")
sh.addShardTag("shard3RS", "ASIA")

# Define zone ranges (for compound key: region + userId)
sh.addTagRange(
  "myDatabase.users",
  { region: "US", userId: MinKey },
  { region: "US", userId: MaxKey },
  "US"
)

sh.addTagRange(
  "myDatabase.users",
  { region: "EU", userId: MinKey },
  { region: "EU", userId: MaxKey },
  "EU"
)

sh.addTagRange(
  "myDatabase.users",
  { region: "ASIA", userId: MinKey },
  { region: "ASIA", userId: MaxKey },
  "ASIA"
)

# Now US users go to shard1, EU to shard2, ASIA to shard3!
# Perfect for GDPR compliance and low latency!

βš–οΈ The Balancer: Automatic Load Distribution

The balancer is a background process that automatically migrates chunks between shards to maintain even data distribution.

🎯 What it Does

  • Monitors chunk distribution
  • Detects imbalanced shards
  • Migrates chunks automatically
  • Runs during off-peak hours
  • Maintains cluster health

βš™οΈ Configuration

// Check balancer status
sh.getBalancerState()

// Stop balancer
sh.stopBalancer()

// Start balancer
sh.startBalancer()

// Schedule balancer window
sh.setBalancerState(true)

🎯 Query Routing: Targeted vs Broadcast

βœ… Targeted Query (Fast)

Query includes shard key. mongos routes to specific shard only.

// Shard key: { userId: 1 }
db.users.find({ userId: 12345 })

// Goes to ONE shard only! ⚑

❌ Broadcast Query (Slow)

Query does NOT include shard key. mongos queries ALL shards.

// Shard key: { userId: 1 }
db.users.find({ email: "..." })

// Goes to ALL shards! 🐌

πŸ’‘ Pro Tip: Always include shard key in queries when possible for maximum performance. Use explain() to verify query routing!

🎯 Sharding Best Practices

πŸ”‘

Choose Shard Key Carefully

It's immutable! Consider cardinality, distribution, and query patterns. Test with production data before committing.

πŸ“Š

Monitor Chunk Distribution

Regularly check balance across shards. Uneven distribution = performance issues and hot spots.

βš–οΈ

Let Balancer Run

Don't disable balancer unless necessary. Schedule balancing windows during low-traffic periods.

🎯

Include Shard Key in Queries

Targeted queries are 10-100x faster than broadcast queries. Always use shard key in WHERE clauses.

πŸ”

Use Replica Sets for Shards

Each shard should be a replica set for high availability. Never use standalone servers as shards.

πŸ“ˆ

Plan for Growth

Add shards before hitting 80% capacity. Pre-splitting chunks helps with initial data loads.

πŸ”

Use explain() Regularly

Verify queries are using shards efficiently. Check for SHARD_MERGE vs full scans.

πŸ’Ύ

Right-Size Chunk Size

Default 128MB works for most. Smaller chunks = more migrations. Larger = less balanced distribution.

πŸ“Š Monitoring Sharded Clusters

Essential Monitoring Commands
# Complete cluster status
sh.status()

# Check shard distribution for a collection
db.users.getShardDistribution()

# Balancer status
sh.getBalancerState()
sh.isBalancerRunning()

# List all shards
db.adminCommand({ listShards: 1 })

# Check chunk counts per shard
db.chunks.aggregate([
  { $group: { _id: "$shard", count: { $sum: 1 } } },
  { $sort: { count: -1 } }
])

# Query performance analysis
db.users.find({ userId: 12345 }).explain("executionStats")

# Config server connection
db.adminCommand({ connPoolStats: 1 })

🚨 Key Metrics to Watch

Chunk Distribution

Shards should have similar chunk counts (Β±20%). Large differences indicate imbalance.

Balancer Activity

Monitor migration frequency. Too many = poor shard key. Too few = possible issues.

Query Performance

Check for SHARD_MERGE in explain(). Broadcast queries = optimization opportunity.

Shard Resources

Monitor CPU, RAM, disk I/O per shard. Hot shards = bottleneck.

❓ Interview Questions & Answers

Q1 What is sharding? When should you use it? β–Ό

Answer:

What is Sharding:

Sharding is MongoDB's method for distributing data across multiple machines (horizontal scaling). It splits data into chunks and distributes them across multiple servers called shards, each of which is an independent MongoDB replica set.

When to Use Sharding:

  • Large datasets: When database size exceeds single server capacity (>1TB)
  • High throughput: When read/write operations exceed single server capacity
  • Working set > RAM: When active dataset doesn't fit in available RAM
  • Geographic distribution: When you need data locality for compliance or latency

When NOT to Use:

  • Small databases (<100GB) - overhead not worth it
  • Low traffic applications - single server sufficient
  • Complex join-heavy workloads - sharding doesn't help with joins
  • Budget constraints - minimum 3 config servers + 2+ shards + mongos = expensive

πŸ—οΈ Sharding Architecture: The 3 Main Components

A MongoDB sharded cluster has three critical components working together:

πŸ’Ύ

1. Shards (Data Storage)

The actual data containers

What it is: Each shard is a separate MongoDB replica set that stores a portion of your data.

Think of it as: Individual storage units in a warehouse. Each unit stores specific items based on a key (like aisle number).

Example: Shard-1 stores users with IDs 0-250,000, Shard-2 stores 250,001-500,000, and so on.

πŸ’‘ Key Point: Each shard is itself a replica set with Primary + Secondary nodes for high availability!

πŸ—ΊοΈ

2. Config Servers (Metadata Brain)

The cluster's navigation system

What it is: Config servers store the metadata and configuration for the entire cluster. They track which data lives on which shard.

What they store:

  • Shard key ranges (chunk mappings)
  • Which shard contains which chunks
  • Cluster topology information
  • Authentication and authorization data

πŸ’‘ Key Point: Config servers are CRITICAL! Without them, mongos can't route queries. Always deploy 3 config servers as a replica set.

🚦

3. Mongos (Query Router)

The traffic controller

What it is: Mongos is a lightweight routing process that directs client queries to the appropriate shard(s). Your application connects to mongos, not directly to shards.

How it works:

  1. Receives query from application
  2. Checks config servers to find which shard has the data
  3. Routes query to correct shard(s)
  4. Merges results and sends back to application

πŸ’‘ Key Point: Mongos has no persistent state! You can run multiple mongos instances for load balancing. If one crashes, use another.

πŸ“ How a Query Flows:

  1. Application sends query: db.users.find({userId: 123456})
  2. Mongos receives query and checks Config Servers for shard key range
  3. Config Servers respond: "User 123456 is in Shard 1"
  4. Mongos routes query to Shard 1 only (not all shards!)
  5. Shard 1 executes query and returns result
  6. Mongos sends result back to Application

πŸ”‘ Shard Keys: The Most Important Decision

⚠️ WARNING: Choosing the wrong shard key can DESTROY your performance. Once set, it's extremely difficult to change! Read this section carefully.

What is a Shard Key?

A shard key is a field (or combination of fields) that MongoDB uses to determine which shard should store each document.

Example: If you shard on userId, MongoDB divides users by their ID: Users 0-250K go to Shard 1, 250K-500K to Shard 2, etc.

🎯 Properties of a Good Shard Key

βœ… 1. High Cardinality

Definition: The field has many unique values (not just a few).

βœ… Good: userId

1 million unique values β†’ Distributes evenly

❌ Bad: gender

Only 2-3 values β†’ All data clumps

βœ… 2. Even Distribution

Definition: Values are evenly distributed, not clustered.

βœ… Good: random userId

Users distributed evenly

❌ Bad: country

80% users in one shard

βœ… 3. Query Isolation

Definition: Most queries include the shard key, so mongos can route to ONE shard.

βœ… Good

find({userId: 123}) β†’ ONE shard

❌ Bad

find({email: "x@y.com"}) β†’ ALL shards

❌ Examples of TERRIBLE Shard Keys

❌ Sharding on "status"

Only 3 values (active, inactive, pending). All data clumps on 3 shards.

❌ Sharding on "createdAt"

All new writes hit the "latest" shard. Creates a write hotspot.

πŸ“¦ Chunks & the Balancer

What are Chunks?

MongoDB groups documents into chunks (contiguous ranges of shard key values). Each chunk is 64MB by default.

Think of it as: Instead of moving individual books, you move entire boxes containing 50 books each. Much more efficient!

βš–οΈ The Balancer

The balancer is a background process that monitors chunk distribution. If one shard has significantly more chunks, it automatically migrates chunks to maintain balance.

πŸ“Š Before Balancing

Shard 1: 80 chunks ⚠️
Shard 2: 40 chunks
Shard 3: 30 chunks
Shard 4: 30 chunks

βœ… After Balancing

Shard 1: 45 chunks βœ…
Shard 2: 45 chunks βœ…
Shard 3: 45 chunks βœ…
Shard 4: 45 chunks βœ…

🎯 Sharding Strategies

1️⃣ Range-Based Sharding

Data is distributed based on ranges of shard key values.

sh.shardCollection("mydb.users", { userId: 1 })

// Shard 1: userId 0 - 250,000
// Shard 2: userId 250,001 - 500,000
// Shard 3: userId 500,001 - 750,000

Use when: Your queries filter by ranges (e.g., dates, sequential IDs)

2️⃣ Hashed Sharding

MongoDB hashes the shard key value and distributes based on the hash.

sh.shardCollection("mydb.users", { userId: "hashed" })

// userId 1 β†’ hash: 8234... β†’ Shard 2
// userId 2 β†’ hash: 1928... β†’ Shard 1
// userId 3 β†’ hash: 9384... β†’ Shard 4

Use when: Shard key is monotonically increasing (timestamps, auto-increment IDs)

3️⃣ Zone Sharding (Geographic)

Assign specific data ranges to specific shards based on location.

// Create zones
sh.addShardToZone("shard-mumbai", "INDIA")
sh.addShardToZone("shard-virginia", "USA")

// Assign ranges to zones
sh.updateZoneKeyRange("mydb.users", 
  { country: "IN" }, { country: "IN" }, "INDIA")

Use when: You need geographic data residency or locality

βš™οΈ Sharding Setup Guide

⚠️ Setting up sharding is complex! Practice in a test environment first.

πŸ“ Step-by-Step Setup

Step 1: Start Config Server Replica Set

# Start 3 config servers
mongod --configsvr --replSet configReplSet --port 27019 --dbpath /data/config1
mongod --configsvr --replSet configReplSet --port 27020 --dbpath /data/config2
mongod --configsvr --replSet configReplSet --port 27021 --dbpath /data/config3

# Initialize replica set
mongosh --port 27019
rs.initiate({
  _id: "configReplSet",
  configsvr: true,
  members: [
    { _id: 0, host: "localhost:27019" },
    { _id: 1, host: "localhost:27020" },
    { _id: 2, host: "localhost:27021" }
  ]
})

Step 2: Start Shard Replica Sets

# Start Shard 1
mongod --shardsvr --replSet shard1 --port 27001 --dbpath /data/shard1
# ... initialize replica set

# Start Shard 2
mongod --shardsvr --replSet shard2 --port 27002 --dbpath /data/shard2
# ... initialize replica set

Step 3: Start Mongos

mongos --configdb configReplSet/localhost:27019,localhost:27020,localhost:27021 --port 27017

Step 4: Add Shards to Cluster

mongosh --port 27017

sh.addShard("shard1/localhost:27001")
sh.addShard("shard2/localhost:27002")

// Verify
sh.status()

Step 5: Enable Sharding & Shard Collection

// Enable sharding on database
sh.enableSharding("mydb")

// Create index on shard key
db.users.createIndex({ userId: 1 })

// Shard the collection
sh.shardCollection("mydb.users", { userId: 1 })

// Verify
db.users.getShardDistribution()

πŸ”§ Sharding Operations

πŸ“Š Monitoring Commands

// Check cluster status
sh.status()

// Check shard distribution
db.collection.getShardDistribution()

// Check balancer status
sh.getBalancerState()

// View chunk information
sh.getBalancerHost()
db.chunks.find().pretty()

⚑ Common Operations

Add a New Shard

sh.addShard("shard3/localhost:27003")

Remove a Shard

db.adminCommand({ removeShard: "shard2" })

Start/Stop Balancer

sh.startBalancer() / sh.stopBalancer()

✨ Sharding Best Practices

βœ… DO's

  • Choose shard key carefully - It's nearly impossible to change later
  • Include shard key in queries - Ensures queries go to one shard
  • Monitor chunk distribution - Watch for imbalances
  • Use hashed sharding - For monotonically increasing keys
  • Test with production-like data - Before going live
  • Deploy config servers as replica sets - For high availability
  • Run multiple mongos instances - For load balancing

❌ DON'Ts

  • Don't shard small collections - Overhead isn't worth it (<100GB)
  • Don't use low-cardinality shard keys - Like "status" or "gender"
  • Don't ignore hotspots - Monitor and address them quickly
  • Don't shard without indexes - Create index on shard key first
  • Don't use timestamps as shard key - Unless hashed
  • Don't manually split chunks - Let MongoDB handle it
  • Don't connect directly to shards - Always use mongos

πŸ’‘ Pro Tip: The Golden Rule

"The best shard key is one that appears in most of your queries, has high cardinality, and distributes data evenly."

🌍 Real-World Examples

πŸ›’ E-commerce Platform (Amazon-scale)

Challenge: 500M products, 1B orders/year, 100M active users

Solution:

  • Shard products by productId (hashed)
  • Shard orders by userId (range-based)
  • 32 shards across 3 geographic regions

Result: 99.99% uptime, <200ms query time, handles Black Friday traffic

πŸ“± Social Media App (Instagram-scale)

Challenge: 2B users, 100M posts/day, real-time feeds

Solution:

  • Shard users by userId (hashed)
  • Shard posts by userId (keeps user's posts together)
  • 64 shards globally distributed

Result: 1M writes/sec, real-time feed generation, <100ms latency worldwide

πŸ“‘ IoT Platform (Smart Home)

Challenge: 100M devices, 10B sensor readings/day, time-series data

Solution:

  • Shard readings by {deviceId: 1, timestamp: 1} (compound)
  • Hashed on deviceId for even distribution
  • 16 shards with time-based data retention

Result: 100K writes/sec, efficient time-range queries, automatic archival

🎯 Interview Questions & Answers

Q2 Explain the architecture of a MongoDB sharded cluster. β–Ό

Answer:

A sharded cluster consists of three main components:

1. mongos (Query Router):

  • Acts as interface between application and cluster
  • Routes queries to appropriate shard(s)
  • Merges results from multiple shards
  • Stateless - can have multiple mongos instances
  • Caches metadata for faster routing

2. Config Servers (Metadata Store):

  • Store cluster metadata (chunk ranges, shard locations)
  • Always deployed as a replica set (3+ servers)
  • Critical - if down, cluster becomes read-only
  • Contain: chunk mappings, zone configurations, balancer state

3. Shards (Data Storage):

  • Each shard is a replica set storing subset of data
  • Contains specific chunks based on shard key ranges
  • Can scale by adding more shards
  • Typically 2-3 members per shard for high availability

Data Flow: Application β†’ mongos β†’ (checks config servers) β†’ routes to shard(s) β†’ returns merged results

Q3 What is a shard key? How do you choose a good one? β–Ό

Answer:

Shard Key: The indexed field(s) MongoDB uses to partition data. It's IMMUTABLE after selection!

Characteristics of Good Shard Keys:

  • High Cardinality: Many unique values (userId, email) not few (gender, boolean)
  • Even Distribution: Values spread evenly, not clustered (hashed keys help)
  • Query Isolation: Most queries should target single shard (include in WHERE)
  • Not Monotonic: Avoid always-increasing values (_id, timestamp) = all writes to one shard
  • Frequently Queried: Should appear in most query filters

Examples:

  • βœ… Good: { userId: 1 }, { email: "hashed" }, { region: 1, userId: 1 }
  • ❌ Bad: { _id: 1 }, { createdAt: 1 }, { status: 1 }

Selection Process:

  1. Analyze query patterns (what fields in WHERE clauses?)
  2. Check cardinality (how many unique values?)
  3. Test distribution (will data spread evenly?)
  4. Consider growth (how will writes distribute over time?)
  5. Prototype with production-like data
Q4 What's the difference between ranged and hashed sharding? β–Ό

Answer:

Ranged Sharding:

  • How: Divides data based on ranges of shard key values
  • Distribution: Documents with nearby values on same shard
  • Pros: Excellent for range queries, sorted access, related data together
  • Cons: Risk of hot spots, uneven distribution with monotonic keys
  • Use when: Range queries common, need data locality
  • Example: userId 0-33M on shard1, 33M-66M on shard2, etc.

Hashed Sharding:

  • How: Uses hash function on shard key for distribution
  • Distribution: Guarantees even distribution regardless of key pattern
  • Pros: Perfect balance, no hot spots, works with monotonic keys
  • Cons: Range queries = broadcast to all shards, can't sort by shard key
  • Use when: Even distribution critical, point queries common
  • Example: hash(userId) distributes randomly across all shards

Choice Depends On:

  • Query patterns: Range queries β†’ Ranged, Point queries β†’ Either
  • Key characteristics: Monotonic β†’ Hashed, Well-distributed β†’ Ranged
  • Hot spots: Concern β†’ Hashed, Not concern β†’ Ranged
Q5 What are chunks? How does chunk splitting and migration work? β–Ό

Answer:

Chunks:

A chunk is a contiguous range of shard key values. MongoDB doesn't distribute individual documents - it distributes chunks.

  • Default size: 128MB (configurable 1MB-1024MB)
  • Indivisible: A chunk lives entirely on one shard
  • Metadata: Chunk ranges stored in config servers

Chunk Splitting:

  • Trigger: When chunk exceeds size limit (128MB)
  • Process: MongoDB automatically splits into two smaller chunks
  • Example: Chunk 0-50M (150MB) β†’ splits to 0-25M (75MB) + 25M-50M (75MB)
  • Location: Both new chunks remain on same shard initially

Chunk Migration:

  • Purpose: Maintain even distribution across shards
  • Trigger: Balancer detects imbalance (shard has too many/few chunks)
  • Process: Balancer moves chunks from overloaded to underloaded shards
  • Timing: Runs in background during configurable window
  • Impact: Minimal - happens incrementally, doesn't block operations

Example Flow:

  1. Chunk grows to 150MB β†’ splits into two 75MB chunks
  2. Shard1 now has 100 chunks, Shard2 has 50
  3. Balancer migrates 25 chunks from Shard1 to Shard2
  4. Result: Both shards have ~75 chunks (balanced)
Q6 What is the balancer and how does it work? β–Ό

Answer:

The Balancer: A background process running on the primary of the config server replica set that maintains even chunk distribution across shards.

How It Works:

  1. Monitors: Continuously checks chunk distribution across shards
  2. Detects Imbalance: Compares chunk counts per shard
  3. Calculates Migrations: Determines which chunks to move where
  4. Migrates: Moves chunks from overloaded to underloaded shards
  5. Updates Metadata: Updates config servers with new chunk locations

Key Characteristics:

  • Automatic: Runs without manual intervention
  • Throttled: Limits concurrent migrations to avoid overload
  • Schedulable: Can configure active windows (e.g., off-peak hours)
  • Pauseable: Can stop/start for maintenance

Configuration:

sh.getBalancerState()      // Check if running
sh.stopBalancer()          // Stop balancer
sh.startBalancer()         // Start balancer
sh.setBalancerState(true)  // Enable balancer

When to Disable:

  • During backups (to ensure consistent snapshots)
  • Large bulk data loads (re-enable after)
  • Maintenance windows on shards
  • Critical high-traffic periods
Q7 Explain targeted vs broadcast queries in sharded clusters. β–Ό

Answer:

Targeted Query:

  • Definition: Query includes shard key in WHERE clause
  • Routing: mongos routes to specific shard(s) only
  • Performance: Very fast - only queries necessary shards
  • Example: db.users.find({ userId: 12345 }) β†’ Goes to ONE shard
  • explain() shows: Single shard accessed

Broadcast Query:

  • Definition: Query does NOT include shard key
  • Routing: mongos sends query to ALL shards
  • Performance: Slow - must query every shard and merge results
  • Example: db.users.find({ email: "..." }) β†’ Goes to ALL shards
  • explain() shows: SHARD_MERGE, all shards queried

Performance Impact:

Targeted Query (1 shard):
  • Latency: 10-50ms
  • Network: Minimal
  • CPU: One shard only
Broadcast Query (10 shards):
  • Latency: 100-500ms (10x worse)
  • Network: 10x traffic
  • CPU: All 10 shards working

Best Practices:

  • Always include shard key when possible
  • Use compound indexes for non-shard-key queries
  • Monitor broadcast query frequency
  • Consider secondary indexes for common non-shard-key queries
Q8 How do you monitor and troubleshoot a sharded cluster? β–Ό

Answer:

Key Monitoring Commands:

sh.status()                        // Overall cluster status
db.collection.getShardDistribution() // Per-collection distribution
sh.getBalancerState()              // Balancer status
db.adminCommand({listShards: 1})   // List all shards

Critical Metrics:

  • Chunk Distribution: Should be roughly equal (Β±20%) across shards
  • Balancer Activity: Check migration frequency and failures
  • Query Performance: Use explain() to verify targeted queries
  • Shard Resources: CPU, RAM, disk I/O per shard
  • Connection Pools: mongos to config servers and shards

Common Issues & Solutions:

  • Uneven Distribution:
    • Check shard key choice (might be poor)
    • Verify balancer is running
    • Consider pre-splitting chunks
  • Slow Queries:
    • Check if broadcast queries (add shard key)
    • Verify indexes on each shard
    • Use explain() to diagnose
  • Hot Shard:
    • Monotonic shard key causing all writes to one shard
    • Solution: Switch to hashed sharding or compound key
  • Balancer Not Running:
    • Check if manually disabled
    • Verify config server primary election
    • Check for failed migrations in logs

Monitoring Tools:

  • MongoDB Atlas (built-in monitoring)
  • MongoDB Ops Manager / Cloud Manager
  • Third-party: Datadog, New Relic, Prometheus
  • Custom scripts using db.serverStatus()
Q1 What is sharding in MongoDB and why is it needed? β–Ό

Answer:

Sharding is MongoDB's horizontal scaling solution that distributes data across multiple servers (shards). It's needed when:

  • Data volume exceeds single server capacity: When your database grows beyond 200-500GB, a single server can't keep the entire working set in RAM
  • Throughput requirements are high: When you need more than 10,000 operations/second, single server becomes a bottleneck
  • Geographic distribution required: To reduce latency for global users by placing data closer to them

Example: If you have 1TB of data and 100,000 queries/sec, sharding across 10 servers means each handles 100GB and 10,000 queries/sec - much more manageable!

Q2 Explain the components of a MongoDB sharded cluster. β–Ό

Answer:

A MongoDB sharded cluster has three essential components:

  1. Shards: Actual database instances (replica sets) that store the data. Each shard holds a portion of the total dataset
  2. Config Servers: Store metadata about the cluster - which chunks are on which shards, cluster topology, etc. Deployed as a replica set of 3 servers
  3. Mongos (Query Router): Lightweight routing process that directs queries to the appropriate shard(s). Applications connect to mongos, not directly to shards

Analogy: Think of shards as warehouses storing products, config servers as the inventory management system that knows what's where, and mongos as the receptionist who directs you to the right warehouse.

Q3 What makes a good shard key? Give examples of good and bad shard keys. β–Ό

Answer:

A good shard key has these properties:

  • High Cardinality: Many unique values (e.g., userId with millions of values)
  • Even Distribution: Values spread evenly across ranges
  • Query Isolation: Most queries include the shard key
  • Non-Monotonic: Avoids values that always increase (unless hashed)

Good Examples:

  • userId (random/hashed) - millions of values, even distribution
  • {country: 1, userId: 1} (compound) - geographic + user distribution
  • orderId (hashed) - prevents hotspots from sequential IDs

Bad Examples:

  • status - only 2-3 values, terrible cardinality
  • createdAt (timestamp) - all new writes hit one shard
  • country - if 80% users in one country, huge imbalance
Q4 Explain chunks and the balancer in MongoDB sharding. β–Ό

Answer:

Chunks: MongoDB groups documents into chunks - contiguous ranges of shard key values. Default chunk size is 64MB.

Example: If sharding on userId:

  • Chunk 1: userId 0 to 100,000 β†’ Shard A
  • Chunk 2: userId 100,001 to 200,000 β†’ Shard B

The Balancer: A background process that monitors chunk distribution. If one shard has significantly more chunks than others (threshold: 8 chunks difference), it automatically migrates chunks to balance the load.

Process:

  1. Balancer checks chunk distribution every few seconds
  2. Identifies imbalanced shards
  3. Moves chunks from overloaded to underloaded shards
  4. Migration happens with minimal impact (uses internal migration protocol)
Q5 What's the difference between range-based and hashed sharding? β–Ό

Answer:

Range-Based Sharding:

  • Distributes data based on ranges of shard key values
  • Example: userId 0-250K β†’ Shard 1, 250K-500K β†’ Shard 2
  • Pros: Efficient for range queries, keeps related data together
  • Cons: Can create hotspots if data isn't evenly distributed
  • Use when: Queries often use ranges (dates, IDs)

Hashed Sharding:

  • MongoDB hashes the shard key value, distributes based on hash
  • Example: userId 1 β†’ hash 8234... β†’ Shard 2
  • Pros: Even distribution guaranteed, prevents hotspots
  • Cons: Range queries hit all shards (scatter-gather)
  • Use when: Shard key is monotonically increasing

Decision: Use hashed for timestamps/auto-increment IDs, use range for naturally distributed data.

Q6 How do queries work in a sharded cluster? Explain targeted vs scatter-gather queries. β–Ό

Answer:

Targeted Query (Best case):

  • Query includes shard key: db.users.find({userId: 12345})
  • Mongos checks config servers, determines exact shard
  • Routes query to ONE specific shard
  • Performance: Fast! Only one shard processes it

Scatter-Gather Query (Worst case):

  • Query doesn't include shard key: db.users.find({email: "user@example.com"})
  • Mongos can't determine which shard has the data
  • Sends query to ALL shards (scatter)
  • Waits for all responses, merges results (gather)
  • Performance: Slow! All shards must be queried

Best Practice: Design your shard key so 80%+ of queries are targeted. If most queries don't include shard key, you chose the wrong shard key!

Q7 When should you NOT use sharding? β–Ό

Answer:

Don't use sharding when:

  1. Data size is small: < 100-200GB. Sharding adds complexity without benefit
  2. Low traffic: < 5,000 ops/sec. Single server can handle it
  3. You can't find a good shard key: If no field has high cardinality and appears in most queries
  4. Complex transactions: If you need ACID across multiple documents frequently (cross-shard transactions are expensive)
  5. Team lacks expertise: Sharding requires DevOps knowledge to maintain

Better alternatives:

  • Vertical scaling (bigger server) for < 200GB
  • Read replicas for read-heavy workloads
  • Archiving old data to reduce active dataset

Rule of Thumb: Shard when you NEED to, not because it's cool technology!

Q8 How would you design sharding for a social media application? β–Ό

Answer (Design approach):

Requirements Analysis:

  • 100M users, 10M posts/day
  • Most queries: "show user's posts", "show user's feed"
  • Need: high write throughput, fast user-specific queries

Shard Key Decision:

  • Users collection: Shard on {userId: "hashed"}
    • Hashed prevents hotspots from sequential user registration
    • Most queries are user-specific: find({userId: X})
  • Posts collection: Shard on {userId: 1, postId: 1} (compound)
    • Keeps all user's posts on same shard (efficient for "my posts" query)
    • postId provides uniqueness

Cluster Setup:

  • 16 shards (can handle 1.6B users at 100M per shard)
  • Geographic distribution: US (8 shards), EU (4 shards), Asia (4 shards)
  • Zone sharding for data residency compliance

Result: Targeted queries (fast!), even distribution, infinitely scalable

Q9 What happens if a shard goes down in a sharded cluster? β–Ό

Answer:

If a shard (replica set) goes down:

  1. Immediate Impact:
    • Queries targeting that shard fail or timeout
    • Queries targeting other shards continue working
    • Only data on the failed shard is unavailable
  2. Replica Set Failover (if properly configured):
    • Primary node fails β†’ Secondary elected as new Primary (10-30 seconds)
    • Queries resume automatically
    • This is why each shard MUST be a replica set!
  3. If entire shard is lost (all replicas):
    • Data on that shard is lost (unless you have backups)
    • Other shards continue functioning
    • You need to restore from backup and re-add the shard

Best Practices for High Availability:

  • Each shard = 3-node replica set (1 Primary, 2 Secondaries)
  • Deploy nodes across different availability zones
  • Regular backups with point-in-time recovery
  • Monitor shard health continuously
Q10 Explain zone sharding and when you would use it. β–Ό

Answer:

Zone Sharding (Geographic Sharding): Assigns specific ranges of shard key values to specific shards based on tags/zones.

How it works:

  1. Tag shards with zones: sh.addShardToZone("shard-mumbai", "INDIA")
  2. Assign data ranges to zones: sh.updateZoneKeyRange("mydb.users", {country: "IN"}, {country: "IN"}, "INDIA")
  3. MongoDB ensures data matching the range stays in that zone

Use Cases:

  • Data Residency Compliance: EU users' data must stay in EU (GDPR)
  • Latency Optimization: Indian users β†’ Mumbai shard, US users β†’ Virginia shard
  • Hardware Optimization: Hot data on SSD shards, cold data on HDD shards

Example Setup:

// Create zones
sh.addShardToZone("shard-us", "USA")
sh.addShardToZone("shard-eu", "EUROPE")
sh.addShardToZone("shard-in", "INDIA")

// Assign ranges
sh.updateZoneKeyRange("mydb.users",
  { country: "US" }, { country: "US" }, "USA")
sh.updateZoneKeyRange("mydb.users",
  { country: "GB" }, { country: "GB" }, "EUROPE")

Result: US users' data stays in US, EU in EU, guaranteed!