Section 8: Distributed Systems

πŸ”„ MongoDB in Distributed Environments

Master how MongoDB scales globally with replica sets, sharding, and multi-region deployments for production-grade distributed applications

πŸ“– The Global Expansion Story: ShopGlobal's Journey

Meet ShopGlobal, an Indian e-commerce startup founded in Mumbai in 2019. What started as a small online marketplace has grown into a platform serving 25 million customers across 40 countries.

In the early days, Arjun, the CTO, ran everything on a single MongoDB server. Life was simple – until their first viral moment during Diwali 2021.

"Our single server crashed at 50,000 concurrent users," Arjun recalls. "We lost β‚Ή2 crore in sales in 4 hours. That's when I knew we needed to go distributed."

Over the next two years, Arjun transformed ShopGlobal's infrastructure:

  • Year 1: Deployed a 3-node replica set for high availability
  • Year 2: Implemented sharding to handle 500K orders/day
  • Year 3: Built a global cluster across Mumbai, Singapore, and Frankfurt

Today, ShopGlobal handles 10 million requests per day with 99.99% uptime!

🌐 What is Distributed MongoDB?

Distributed MongoDB runs across multiple servers, data centers, or regions to achieve high availability, scalability, and low latency for users worldwide.

Simple Analogy

Think of distributed MongoDB like a chain of restaurants. Instead of one giant restaurant that gets overwhelmed during rush hour, you have multiple locations (servers) that share the load. If one closes temporarily, customers go to another seamlessly.

πŸ”„
Replica Sets
Multiple copies of your data. If one server fails, another takes over automatically. High availability.
🧩
Sharding
Splits data across multiple servers (shards). Each holds a portion, enabling horizontal scaling.
🌍
Global Clusters
Deploy across regions. Users connect to nearest data center for low latency.
πŸ”€
mongos Routers
Query routers that direct requests to the right shard automatically.

πŸ”„ Replica Sets: Your First Line of Defense

A replica set is a group of MongoDB servers maintaining the same data. It's the foundation of high availability.

MongoDB Replica Set Architecture Application / Driver PRIMARYReads + Writes SECONDARYReads (optional) SECONDARYReads (optional) ⟢ Replication (Oplog)
ComponentRoleDetails
PrimaryHandles all writesOnly one primary. Replicates to secondaries.
SecondaryData copiesMaintains copies. Can serve reads. Becomes primary during failover.
ArbiterVoting onlyParticipates in elections but holds no data.
OplogReplication logCapped collection recording all write operations.
MongoDB Shell
// Initialize a replica set
rs.initiate({
  _id: "myReplicaSet",
  members: [
    { _id: 0, host: "server1:27017", priority: 2 },
    { _id: 1, host: "server2:27018", priority: 1 },
    { _id: 2, host: "server3:27019", priority: 1 }
  ]
})

// Check status
rs.status()
Automatic Failover (10-30 seconds)
  1. Secondaries detect primary unavailable (heartbeat timeout: 10s)
  2. Eligible secondary calls for election
  3. Members vote based on priority, oplog position
  4. Candidate with majority votes becomes new primary
  5. Applications automatically reconnect

🧩 Sharding: Scaling to Billions of Documents

Sharding distributes data across multiple servers, allowing you to handle datasets larger than any single server can hold.

When Do You Need Sharding?
  • Single server can't handle write throughput
  • Dataset exceeds single server storage
  • Working set exceeds RAM
  • Need geographic data distribution
MongoDB Sharded Cluster Application mongos Router mongos Router mongos Router Config Server Replica Set Shard 1 (Replica Set) PRIMARY SECONDARY SECONDARY Shard 2 (Replica Set) PRIMARY SECONDARY SECONDARY Shard 3 (Replica Set) PRIMARY SECONDARY SECONDARY

Shard Key Selection

MongoDB Shell
// Enable sharding on database
sh.enableSharding("ecommerce")

// Create index on shard key
db.orders.createIndex({ customerId: 1 })

// Shard the collection (hashed for even distribution)
sh.shardCollection(
  "ecommerce.orders",
  { customerId: "hashed" }
)

// Check status
sh.status()
CharacteristicGood Shard KeyBad Shard Key
CardinalityHigh (many unique values) - userIdLow - status, country
DistributionEven - hashed(userId)Uneven (hot spots) - createdAt
Query PatternsFrequently queried - customerIdRarely used in queries
MonotonicityNon-monotonic or hashedAlways increasing - _id, timestamp

🌍 Global Clusters: Worldwide Deployment

Deploy MongoDB across multiple geographic regions for low latency and data residency compliance.

⚑
Low Latency
Users connect to nearest region. Mumbai users hit Mumbai servers (5ms), not Frankfurt (150ms).
πŸ“œ
Data Residency
Keep EU data in EU, Indian data in India. GDPR compliance.
πŸ›‘οΈ
Disaster Recovery
If entire region goes down, traffic routes to another region.
πŸ”„
Zone Sharding
Control which data lives where by country, tier, or any field.
MongoDB Shell
// Create zones for different regions
sh.addShardTag("shard-mumbai", "INDIA")
sh.addShardTag("shard-frankfurt", "EU")

// Define zone ranges - Indian data stays in Mumbai
sh.addTagRange(
  "ecommerce.customers",
  { region: "IN" },
  { region: "IN" + MaxKey },
  "INDIA"
)

πŸ” Consistency Models

Write Concern

MongoDB Shell
// WEAK - Fire and forget
db.orders.insertOne(doc, { writeConcern: { w: 0 } })

// DEFAULT - Primary only
db.orders.insertOne(doc, { writeConcern: { w: 1 } })

// STRONG - Majority (recommended)
db.orders.insertOne(doc, { 
  writeConcern: { w: "majority", j: true } 
})

Read Preference

PreferenceBehaviorUse Case
primaryAlways from primaryStrong consistency
primaryPreferredPrimary if availablePrefer consistency
secondaryOnly secondariesRead scaling, analytics
secondaryPreferredSecondary if availableOffload primary
nearestLowest latency nodeGeo-distributed apps

πŸš€ Deployment Strategies

🏠
Single Data Center

All nodes in one DC. Simple but single point of failure.

  • 3-node replica set
  • Low latency between nodes
  • Risk: DC failure = downtime
  • Best for: Development, non-critical
🏒
Multi Data Center (Recommended)

Nodes spread across 2-3 DCs. Survives single DC failure.

  • Primary + Secondary in DC1
  • Secondary in DC2, Arbiter in DC3
  • Survives single DC failure
  • Best for: Production
🌐
Global Multi-Region

Sharded cluster across continents with zone sharding.

  • Shards in Mumbai, Singapore, Frankfurt
  • Zone sharding for data residency
  • Local reads everywhere
  • Best for: Global apps

⚠️ Challenges & Solutions

Network Partitions

Problem: Network failure splits cluster.

Solution: Deploy across 3+ availability zones. Only majority partition elects primary.

Replication Lag

Problem: Secondaries fall behind, stale reads.

Solution: Use readConcern: "majority". Monitor with rs.printReplicationInfo().

Shard Hotspots

Problem: One shard gets all traffic.

Solution: Use hashed shard keys. Avoid monotonic keys.

Cross-Shard Queries

Problem: Queries without shard key hit all shards.

Solution: Include shard key in queries. Design model to avoid cross-shard ops.

βœ… Best Practices

πŸ”’
Use Odd Number of Nodes
Deploy 3, 5, or 7 nodes. Even numbers cause split-brain issues.
πŸ“Š
Monitor Replication Lag
Alert on lag > 10 seconds. Check rs.printSlaveReplicationInfo().
πŸ”‘
Choose Shard Key Carefully
Shard key is immutable! Test with production-like data first.
πŸ’Ύ
Enable Journaling
Always enable (default). Ensures durability during crashes.
πŸ”
Use w: "majority"
For critical writes: writeConcern: { w: "majority", j: true }.
πŸ“ˆ
Plan for Growth
Shard early if expecting growth. Sharding later is painful.

🎯 Interview Questions (25)

1What is a MongoDB replica set?β–Ό

Answer: A replica set is a group of MongoDB instances maintaining the same data. One primary handles writes, secondaries maintain copies.

  • High Availability: Auto failover (10-30s)
  • Data Redundancy: Multiple copies
  • Read Scaling: Secondaries can serve reads
2Explain the election process in MongoDB.β–Ό

Answer: When primary fails:

  1. Secondaries detect unavailability (10s heartbeat)
  2. Eligible secondary calls election
  3. Members vote by priority, oplog position
  4. Majority wins (N/2 + 1 votes)
  5. New primary announced

Takes 10-30 seconds. Writes blocked, reads continue from secondaries.

3What is the oplog?β–Ό

Answer: Operations log - a capped collection at local.oplog.rs recording all writes. Secondaries replay oplog to sync.

  • Idempotent operations
  • Default 5% of disk
  • If secondary falls too far behind, needs full resync
4What is sharding and when to use it?β–Ό

Answer: Horizontal scaling - distributes data across multiple shards.

Use when:

  • Dataset exceeds single server
  • Write throughput exceeds capacity
  • Working set exceeds RAM
  • Need geographic distribution
5What makes a good shard key?β–Ό

Good: High cardinality, even distribution, query isolation, non-monotonic. Example: hashed(userId)

Bad: status (low cardinality), country (uneven), createdAt (monotonic), _id (ObjectId is monotonic)

6Explain write concern.β–Ό
  • w:0 - Fire and forget
  • w:1 - Primary only (default)
  • w:"majority" - Majority of nodes
  • j:true - Written to journal

Higher = stronger durability = slower writes.

7What is read preference?β–Ό
  • primary - Always primary (strong consistency)
  • primaryPreferred - Primary if available
  • secondary - Only secondaries (read scaling)
  • secondaryPreferred - Secondary if available
  • nearest - Lowest latency (geo-distributed)
8What is mongos?β–Ό

Query router in sharded cluster. Routes queries to correct shard(s), merges results, caches metadata. Stateless - scale horizontally.

9What are config servers?β–Ό

Store metadata: chunk locations, cluster config, namespaces. Must be replica set. If unavailable, cluster is read-only.

10Hashed vs ranged sharding?β–Ό

Ranged: Divides by key value ranges. Good for range queries. Risk of uneven distribution.

Hashed: Hash function on key. Even distribution. Poor for range queries.

11What is zone sharding?β–Ό

Associates shard key ranges with specific shards for data locality.

  • Data residency (GDPR)
  • Low latency
  • Tiered storage
12What happens during network partition?β–Ό

Only partition with majority votes can elect primary. Minority becomes read-only. Deploy across 3+ AZs to survive single zone failure.

13What is an arbiter?β–Ό

Votes in elections but holds no data. Use when need odd vote count without full replica cost. Avoid if you can afford 3 full replicas.

14What is replication lag?β–Ό

Delay between primary write and secondary replication. Causes: network latency, slow disk, high load.

Solution: readConcern:"majority", monitor with rs.printReplicationInfo().

15Explain causal consistency.β–Ό

Guarantees causally related operations seen in same order. Solves "read your writes" problem with secondary reads.

const session = client.startSession({ causalConsistency: true });
16What is chunk migration?β–Ό

Moving data chunks between shards to balance cluster. Balancer monitors distribution, initiates migrations. Schedule for low-traffic periods.

17How to handle schema changes in distributed MongoDB?β–Ό
  • Additive: Add fields, app handles missing
  • Lazy migration: Update on access
  • Background migration: Batch updates
  • Index changes: Background creation
18What is WiredTiger?β–Ό

Default storage engine since v3.2. Features: document-level locking, compression (Snappy/Zstd), checkpoints (60s), journal, configurable cache.

19How to backup a distributed cluster?β–Ό
  • mongodump with --oplog
  • Filesystem snapshots (stop balancer)
  • MongoDB Atlas Backup
  • Ops Manager

Stop balancer, backup config servers with shards.

20What is scatter-gather query?β–Ό

Query without shard key hits ALL shards. Avoid by including shard key in queries, design model for query isolation.

21How do transactions work in distributed MongoDB?β–Ό

Multi-doc ACID since v4.0 (replica), v4.2 (sharded). Two-phase commit across shards. 60s timeout. Use sparingly for performance.

22Primary vs secondary reads?β–Ό

Primary: Latest data (strong consistency), adds load.

Secondary: May be stale, distributes load, good for analytics.

23How to monitor distributed MongoDB?β–Ό
  • Replica: rs.status(), replication lag
  • Sharding: sh.status(), chunk distribution
  • Performance: db.serverStatus()
  • Tools: Atlas, Ops Manager, Prometheus+Grafana
24What is retryable writes?β–Ό

Drivers auto-retry failed writes (network errors, elections). Uses unique ID to detect duplicates. Default since v4.2. Only with replica sets.

25How to design multi-tenant with sharding?β–Ό

Option 1: { tenantId: 1, _id: 1 } - query isolation, risk of hot shards

Option 2: { tenantId: "hashed" } - even distribution, loses locality

Option 3: Zone sharding - premium tenants on dedicated shards

Best: Compound key { tenantId: 1, date: 1 }