Section 13: Scenario Based Question

🌐 Distributed Systems Interview Questions

MongoDB Sharding, Replica Sets, and Scaling Architecture

🎯 About Distributed Systems Questions

Distributed systems questions test your understanding of MongoDB's architecture for handling high availability, scalability, and fault tolerance. These are crucial for senior roles and companies dealing with large-scale data.

πŸ’‘ What Interviewers Look For:
  • Architectural Understanding: How MongoDB distributes data across servers
  • Failure Handling: What happens when nodes fail
  • Trade-offs: CAP theorem implications and consistency vs availability
  • Performance: How to scale reads and writes
  • Real-world Experience: Practical deployment considerations
30
Questions
60%
Asked at FAANG
Senior+
Level Focus
⚠️ Common Mistakes to Avoid:
  • Confusing sharding with replication
  • Not understanding election processes
  • Ignoring write concern and read preference implications
  • Oversimplifying CAP theorem trade-offs
  • Not considering network partitions

πŸ”„ Replica Sets

Questions about MongoDB's replication mechanism for high availability.

Q1 Medium Amazon Microsoft
Explain how MongoDB replica sets work. What are the different node types and their roles?

βœ“ Complete Answer:

What is a Replica Set?

A replica set is a group of MongoDB servers that maintain the same dataset, providing redundancy and high availability through data replication.

    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
    β”‚   Primary   β”‚ ◄─── Accepts all writes
    β”‚   (Main)    β”‚ ──── Replicates to secondaries
    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
           β”‚
           β”‚ Replication
           β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
           β–Ό              β–Ό              β–Ό
    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
    β”‚Secondary 1β”‚  β”‚Secondary 2β”‚  β”‚Secondary 3β”‚
    β”‚ (Replica) β”‚  β”‚ (Replica) β”‚  β”‚ (Replica) β”‚
    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                        β”‚
                                        β”‚ Optional
                                        β–Ό
                                  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                                  β”‚  Arbiter  β”‚
                                  β”‚(No data)  β”‚
                                  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Node Types:

1. Primary Node
  • Receives all write operations
  • Records all changes to oplog (operations log)
  • By default, handles read operations
  • Only one primary per replica set
2. Secondary Nodes
  • Replicate data from primary's oplog
  • Can serve read operations (with read preference)
  • Can become primary during election
  • Typically have 2+ secondaries for redundancy
3. Arbiter (Optional)
  • Participates in elections but holds NO data
  • Used for odd number of voting members (avoid ties)
  • Lightweight, minimal resources
  • Cannot become primary

Special Secondary Types:

  • Priority 0: Cannot become primary (for reporting servers)
  • Hidden: Not visible to applications, cannot become primary
  • Delayed: Maintains historical data (e.g., 1 hour behind) for recovery
// Configure a replica set
rs.initiate({
  _id: "myReplicaSet",
  members: [
    { _id: 0, host: "server1:27017", priority: 2 },      // Primary candidate
    { _id: 1, host: "server2:27017", priority: 1 },      // Secondary
    { _id: 2, host: "server3:27017", priority: 0 },      // Priority 0 (won't be primary)
    { _id: 3, host: "server4:27017", arbiterOnly: true } // Arbiter
  ]
})
Q2 Hard Google Meta
Walk me through what happens when a primary node fails in a replica set. How does MongoDB handle the election process?

βœ“ Complete Answer:

Election Process Step-by-Step:

Step 1: Failure Detection (10 seconds default)
  • Secondary nodes send heartbeats every 2 seconds
  • If primary doesn't respond within 10 seconds β†’ failure detected
  • Secondaries initiate election process
Step 2: Election Initiation
  • Eligible secondary calls for an election
  • Must have priority > 0 and not be hidden/delayed
  • Must be up-to-date (within 10 seconds of most recent oplog entry)
Step 3: Voting
  • Each voting member gets one vote
  • Requires majority (n/2 + 1) to win
  • Higher priority nodes more likely to win
  • More up-to-date nodes preferred
Step 4: New Primary Elected
  • Winning secondary becomes primary
  • Starts accepting writes immediately
  • Other secondaries start replicating from new primary
  • Typical election time: 12-30 seconds
Timeline of Election:

t=0s:    Primary fails
t=10s:   Secondaries detect failure (heartbeat timeout)
t=11s:   Election initiated
t=12s:   Votes collected
t=13s:   New primary elected
t=13s+:  New primary accepts writes

During Election (10-13s):
- NO writes accepted (write operations queue or fail)
- Reads still work on secondaries (if configured)
                            

Critical Concepts:

Majority Requirement

Why majority matters: Prevents split-brain scenarios

  • 3-member set: needs 2 votes (can tolerate 1 failure)
  • 5-member set: needs 3 votes (can tolerate 2 failures)
  • 7-member set: needs 4 votes (can tolerate 3 failures)
⚠️ Important Edge Cases:
  • Network Partition: If primary can't reach majority, it steps down
  • No Majority: If
  • Tie Breaker: Use arbiter to ensure odd number of voters
  • Old Primary Returns: Automatically becomes secondary
// Check replica set status during election
rs.status()

// Output shows:
{
  "members": [
    {
      "name": "server1:27017",
      "state": 1,        // PRIMARY
      "stateStr": "PRIMARY",
      "health": 1
    },
    {
      "name": "server2:27017",
      "state": 2,        // SECONDARY
      "stateStr": "SECONDARY",
      "health": 1
    },
    // During election, state might be:
    // state: 0 (STARTUP)
    // state: 6 (UNKNOWN)
    // state: 3 (RECOVERING)
  ]
}
Q3 Medium Netflix
Explain write concern and read preference. How do they affect consistency and performance?

βœ“ Complete Answer:

Write Concern:

Write concern specifies the level of acknowledgment requested from MongoDB for write operations.

Write Concern Levels:
  • w: 0 - No acknowledgment (fire and forget)
  • w: 1 - Acknowledged by primary only (default)
  • w: "majority" - Acknowledged by majority of replica set
  • w: N - Acknowledged by N nodes
// Write concern examples
db.orders.insertOne(
  { customer: "Alice", total: 150 },
  { writeConcern: { w: 1 } }  // Ack from primary only (fast)
)

db.orders.insertOne(
  { customer: "Bob", total: 500 },
  { writeConcern: { w: "majority", wtimeout: 5000 } }  // Safer, slower
)

// With journal option
db.orders.insertOne(
  { customer: "Charlie", total: 1000 },
  { writeConcern: { w: "majority", j: true } }  // Persisted to disk
)

Read Preference:

Read preference describes how MongoDB clients route read operations to members of a replica set.

Read Preference Modes:
  • primary - All reads from primary (default, strongest consistency)
  • primaryPreferred - Primary if available, else secondary
  • secondary - All reads from secondaries (reduces primary load)
  • secondaryPreferred - Secondary if available, else primary
  • nearest - Lowest network latency (any node)
// Read preference examples
db.orders.find({ status: "pending" })
  .readPref("secondary")  // Read from secondary

db.orders.find({ status: "completed" })
  .readPref("primaryPreferred")  // Primary first, fallback to secondary

Trade-offs:

Configuration Consistency Performance Use Case
w:1 + primary ⭐⭐⭐ ⭐⭐⭐ Default, balanced
w:majority + primary ⭐⭐⭐⭐⭐ ⭐⭐ Critical data
w:1 + secondary ⭐⭐ ⭐⭐⭐⭐⭐ Analytics, reports
w:0 + nearest ⭐ ⭐⭐⭐⭐⭐ Logs, metrics
⚠️ Common Pitfall - Read Your Own Writes:

If you write with w:1 to primary, then immediately read from secondary, you might not see your write yet (replication lag).

Solution: Use read preference "primary" for read-after-write consistency, or w:"majority" for critical writes.

πŸ—‚οΈ Sharding

Questions about MongoDB's horizontal scaling mechanism.

Q4 Hard Amazon Uber
Explain MongoDB sharding architecture. What are the components and how does data distribution work?

βœ“ Complete Answer:

What is Sharding?

Sharding is MongoDB's approach to horizontal scaling - distributing data across multiple servers (shards) to handle large datasets and high throughput.

Application
     β”‚
     β–Ό
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚   mongos    β”‚ ◄── Query Router (lightweight, stateless)
β”‚  (Router)   β”‚     Routes queries to correct shards
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
     β”‚
     β”‚ Asks: "Which shards have this data?"
     β–Ό
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚Config Serverβ”‚ ◄── Metadata store (what data is where)
β”‚ Replica Set β”‚     Stores: chunk ranges, shard locations
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
     β”‚
     β”‚ Routes queries based on metadata
     β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
     β–Ό            β–Ό            β–Ό            β–Ό
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚ Shard 1 β”‚  β”‚ Shard 2 β”‚  β”‚ Shard 3 β”‚  β”‚ Shard 4 β”‚
β”‚(Replica β”‚  β”‚(Replica β”‚  β”‚(Replica β”‚  β”‚(Replica β”‚
β”‚  Set)   β”‚  β”‚  Set)   β”‚  β”‚  Set)   β”‚  β”‚  Set)   β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
users: 1-250K  251K-500K   501K-750K   751K-1M

Components:

1. Mongos (Query Router)
  • Application connects to mongos (not directly to shards)
  • Routes queries to appropriate shard(s)
  • Merges results from multiple shards
  • Stateless - can run multiple instances
2. Config Servers
  • Store cluster metadata (which chunks are on which shards)
  • Must be deployed as replica set (3 members minimum)
  • Critical - if unavailable, sharding stops working
  • Small data volume but high importance
3. Shards
  • Each shard is a replica set holding subset of data
  • Minimum 2 shards required for sharded cluster
  • Each shard independent - can scale independently
  • Should be geographically distributed for disaster recovery

How Data Distribution Works:

Chunks

Data divided into contiguous ranges called "chunks" (default 128MB)

  • Each chunk contains documents with shard key in specific range
  • Example: userId 1-10000 in chunk 1, 10001-20000 in chunk 2
  • Chunks automatically split when they exceed size limit
  • Balancer migrates chunks to maintain even distribution
// Enable sharding on database
sh.enableSharding("myDatabase")

// Shard a collection on userId
sh.shardCollection(
  "myDatabase.users",
  { userId: 1 }  // Shard key
)

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

// Output:
Shard shard1 at shard1/server1:27017
  data: 45.2MB docs: 250000 chunks: 4
  
Shard shard2 at shard2/server2:27017
  data: 47.8MB docs: 260000 chunks: 4

Query Routing Logic:

Targeted Query (Best Case)

Query includes shard key β†’ mongos routes to specific shard

db.users.find({ userId: 12345 })  // Goes to 1 shard only
Broadcast Query (Worst Case)

Query doesn't include shard key β†’ mongos queries all shards

db.users.find({ email: "[email protected]" })  // Queries all shards
⚠️ Critical Design Decision - Shard Key Selection:

Cannot be changed after sharding! Choose carefully:

  • High cardinality (many unique values)
  • Good write distribution (avoid hotspots)
  • Included in most queries (for targeted routing)
Q5 Hard LinkedIn Uber
What makes a good shard key? Explain with examples of good and bad shard keys.

βœ“ Complete Answer:

Three Golden Rules for Shard Keys:

1. High Cardinality

Many unique values allow data to be distributed across many chunks

  • Good: userId (millions of unique values)
  • Bad: country (only ~200 values) - creates large undivisible chunks
2. Even Distribution

Values should distribute writes evenly across shards (avoid hotspots)

  • Good: hashed userId
  • Bad: timestamp (all new writes go to same shard holding latest range)
3. Query Isolation

Shard key should appear in most queries for targeted routing

  • Good: userId for user-centric queries
  • Bad: Random field not used in queries (forces broadcast)

Examples:

❌ Bad: Sequential _id or Timestamp
sh.shardCollection("db.orders", { _id: 1 })
// or
sh.shardCollection("db.orders", { createdAt: 1 })

Problem: "Monotonically Increasing" shard key
- All new inserts go to the SAME shard (highest range)
- Creates write hotspot on one shard
- Other shards sit idle
- Severe bottleneck!

Visual:
Time β†’   All writes
Shard 1: [old data, idle]
Shard 2: [old data, idle]
Shard 3: [old data, idle]
Shard 4: [recent data] ◄─── πŸ”₯ 100% of writes here!
βœ… Good: Hashed Shard Key
sh.shardCollection("db.orders", { _id: "hashed" })

Benefits:
- Hash function evenly distributes values
- Writes distributed across all shards
- No hotspots

Drawback:
- Cannot do range queries on shard key
- find({ _id: { $gte: X, $lte: Y } }) becomes broadcast query
βœ… Good: Compound Shard Key
sh.shardCollection("db.orders", { customerId: 1, orderId: 1 })

Benefits:
- customerId provides distribution
- orderId adds granularity
- Queries by customerId are targeted
- Supports range queries within customer

Best for: Multi-tenant applications, user-partitioned data
❌ Bad: Low Cardinality
sh.shardCollection("db.users", { status: 1 })  // "active" or "inactive"

Problem:
- Only 2 possible values
- Can only create 2 chunks maximum
- Cannot distribute beyond 2 shards
- Defeats purpose of sharding!

Real-World Patterns:

Use Case Good Shard Key Why
User profiles { userId: "hashed" } Even distribution, queries by user
IoT sensor data { deviceId: 1, timestamp: 1 } Device isolation + time range queries
Log entries { _id: "hashed" } No hotspots on sequential inserts
E-commerce orders { customerId: 1, orderDate: 1 } Customer isolation + date ranges

βš–οΈ Consistency & CAP Theorem

Questions about data consistency and distributed systems trade-offs.

Q6 Hard Google Meta
Explain the CAP theorem and where MongoDB sits on the spectrum. Can you achieve strong consistency in MongoDB?

βœ“ Complete Answer:

CAP Theorem Basics:

In a distributed system, you can only guarantee 2 out of 3 properties:

  • C (Consistency): All nodes see the same data at the same time
  • A (Availability): Every request gets a response (success or failure)
  • P (Partition Tolerance): System continues despite network failures
CAP Theorem Triangle:

            Consistency
               /\
              /  \
             /    \
            /  CP  \
           /________\
          /    MongoDB  \
         /    (tunable)  \
        /___________________\
Availability              Partition
                         Tolerance

MongoDB: Prioritizes CP, but configurable toward AP

MongoDB's Position:

Default Behavior (CP)

MongoDB prioritizes Consistency and Partition tolerance:

  • During network partition, minority side becomes read-only
  • Primary steps down if it can't reach majority
  • Sacrifices availability to maintain consistency
Strong Consistency Configuration

YES - MongoDB can provide strong consistency:

// Strong consistency: Read your own writes guaranteed
db.collection.insertOne(
  { data: "important" },
  { writeConcern: { w: "majority" } }  // Wait for majority ack
)

db.collection.findOne(
  { data: "important" },
  { readPreference: "primary" }  // Read from primary only
)

Result: If write succeeds, subsequent read WILL see it
Eventual Consistency Configuration

Can tune toward availability (AP):

// Eventual consistency: Prioritize availability
db.collection.insertOne(
  { data: "log_entry" },
  { writeConcern: { w: 1 } }  // Primary only
)

db.collection.findOne(
  { data: "log_entry" },
  { readPreference: "secondaryPreferred" }  // Can read from secondary
)

Result: May read stale data, but always available

Practical Scenarios:

Scenario 1: Network Partition
Setup: 5-node replica set across 2 data centers
- DC1: Primary + 2 Secondaries (3 nodes)
- DC2: 2 Secondaries (2 nodes)

Network partition separates DC1 from DC2

MongoDB Behavior:
βœ“ DC1 (has majority): Continues operating, primary stays up
βœ— DC2 (no majority): All nodes become read-only

Trade-off: DC2 sacrifices availability to prevent split-brain
Scenario 2: Replication Lag
Setup: Write with w:1, read from secondary

1. Write to primary (acknowledged immediately)
2. Replication to secondary (takes 100ms)
3. Read from secondary before replication completes

Result: Read doesn't see the write (stale read)

Prevention: Use w:"majority" + readPreference:"primary"
🎯 Interview Tip:

Don't just say "MongoDB is CP." Explain that it's tunable based on write concern and read preference. Show you understand the trade-offs for different application requirements.

Q7 Medium Amazon
What is the oplog in MongoDB? How does it work and why is it important?

βœ“ Complete Answer:

What is the Oplog?

The oplog (operations log) is a special capped collection that records all write operations that modify data in the primary node. It's the foundation of MongoDB replication.

Key Characteristics:
  • Special capped collection: local.oplog.rs
  • Fixed size (configurable, default ~5% of disk)
  • FIFO - oldest entries automatically removed when full
  • Idempotent operations - can be applied multiple times safely
// View oplog entries
use local
db.oplog.rs.find().sort({ $natural: -1 }).limit(5)

// Sample oplog entry
{
  "ts": Timestamp(1234567890, 1),  // Timestamp
  "t": Long("1"),                   // Term number
  "h": Long("123456789"),           // Hash
  "v": 2,                           // Oplog version
  "op": "i",                        // Operation: i(insert), u(update), d(delete)
  "ns": "mydb.users",               // Namespace (database.collection)
  "o": {                            // Operation document
    "_id": ObjectId("..."),
    "name": "Alice",
    "age": 25
  }
}

How Replication Works via Oplog:

1. Primary receives write
           β”‚
           β–Ό
2. Write applied to primary data
           β”‚
           β–Ό
3. Operation recorded in primary's oplog
           β”‚
           β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
           β–Ό            β–Ό            β–Ό
4. Secondaries tail oplog (continuous polling)
           β”‚            β”‚            β”‚
           β–Ό            β–Ό            β–Ό
5. Secondaries apply operations from oplog
           β”‚            β”‚            β”‚
           β–Ό            β–Ό            β–Ό
6. Secondaries' data matches primary
Why Oplog Size Matters:
  • Too Small: Secondary might not catch up before old entries are overwritten
  • Replication Window: Time it takes to fill the oplog
  • Recovery Time: Determines how long a secondary can be offline
// Check oplog status
rs.printReplicationInfo()

// Output:
configured oplog size:   5120MB
log length start to end: 86400secs (24hrs)
oplog first event time:  Mon Dec 21 2024 00:00:00
oplog last event time:   Tue Dec 22 2024 00:00:00
now:                     Tue Dec 22 2024 06:00:00
⚠️ Critical Scenario - Oplog Too Small:

If secondary is offline for longer than oplog window, it becomes "too stale" and requires full resync (copying entire dataset). Plan oplog size for longest expected downtime.

Q8 Hard Netflix Uber
What is "rollback" in MongoDB? When does it occur and how is it handled?

βœ“ Complete Answer:

What is Rollback?

Rollback occurs when a former primary rejoins the replica set and must undo operations that were not replicated to the new primary before the old primary went offline.

Rollback Scenario:

Time 0: Normal operation
   Primary A ──> Secondary B ──> Secondary C
   [ops 1-100]   [ops 1-100]    [ops 1-100]

Time 1: Network partition (A isolated)
   Primary A (alone)     |   Secondary B + C
   [ops 101-105]         |   [ops 1-100]
                         |
Time 2: B becomes new primary (majority with C)
   (isolated) A          |   Primary B ──> Secondary C
   [ops 101-105]         |   [ops 101-103 different]

Time 3: A reconnects and discovers conflict
   A has ops 101-105 (not on B/C)
   B has ops 101-103 (different operations)
   
   β†’ A must ROLLBACK ops 101-105
   β†’ Then sync from new primary B
When Rollback Happens:
  • Primary gets isolated (can't reach majority)
  • Primary continues accepting writes (from clients before timeout)
  • New primary elected in majority partition
  • Old primary rejoins with divergent operations
Rollback Process:
  1. Old primary detects it's behind new primary
  2. Finds common point in oplog
  3. Saves rolled-back operations to rollback files
  4. Reverts to common point
  5. Syncs from new primary
  6. Becomes secondary in current state
// Rollback files location
/data/db/rollback/

// Example: removed.2024-12-22T06-00-00.0.bson
// Contains documents that were rolled back

// MongoDB 4.0+ prevents rollbacks with w:"majority"
db.criticalData.insertOne(
  { data: "important" },
  { writeConcern: { w: "majority" } }
)
// This write won't be acknowledged until replicated to majority
// Therefore, it CAN'T be rolled back
⚠️ Preventing Rollbacks:
  • Use w:"majority": Ensures writes replicated before ack
  • Monitor replication lag: Alert if secondaries fall behind
  • Size oplog appropriately: Allow time for recovery
  • Handle rollback files: Review and reapply if needed
πŸ’‘ Interview Insight:

Rollback is MongoDB's mechanism to maintain consistency across the replica set. It's a safety feature, not a bug. Show you understand it's a consequence of choosing availability during partitions, and explain mitigation strategies.

⚑ Performance & Scaling

Questions about optimizing and scaling MongoDB deployments.

Q9 Medium LinkedIn
How do you scale MongoDB reads? What strategies can you use?

βœ“ Complete Answer:

Read Scaling Strategies:

1. Read from Secondaries

Distribute read load across replica set members

// Scale reads across secondaries
db.users.find({ city: "NYC" })
  .readPref("secondary")

// Distribute to nearest node
db.users.find({ status: "active" })
  .readPref("nearest")

Benefits:
βœ“ Offloads read traffic from primary
βœ“ Reduces primary CPU/memory pressure
βœ“ Can handle 5x-10x more reads

Trade-off:
βœ— Eventual consistency (may read stale data)
βœ— Replication lag (typically 0-500ms)
2. Add More Secondaries

Horizontal scaling of read capacity

3-node replica set:  1 primary + 2 secondaries
                     β†’ 2x read capacity

7-node replica set:  1 primary + 6 secondaries  
                     β†’ 6x read capacity

Note: Up to 7 voting members max
      Can have up to 50 total members (43 non-voting)
3. Create Read-Only Secondaries

Dedicated secondaries for analytics/reporting

// Configure priority 0 secondary (won't become primary)
rs.reconfig({
  members: [
    { _id: 0, host: "primary:27017", priority: 2 },
    { _id: 1, host: "secondary1:27017", priority: 1 },
    { _id: 2, host: "analytics:27017", priority: 0, hidden: true }
  ]
})

// Analytics queries go here
// Doesn't affect production traffic
4. Shard the Collection

Distribute data across multiple shards for massive scale

No sharding: 1 replica set = 6 secondaries = 6x read capacity

Sharding:    4 shards Γ— 6 secondaries each = 24x read capacity
             
Each shard handles subset of data
Queries can be parallelized across shards
5. Proper Indexing

Most impactful optimization - essential before scaling

// Without index: Collection scan (slow)
db.users.find({ email: "user@example.com" })  // 1000ms

// With index: Index scan (fast)
db.users.createIndex({ email: 1 })
db.users.find({ email: "user@example.com" })  // 5ms

// 200x faster!
// Always optimize queries before adding hardware

Decision Matrix:

Requirement Best Strategy Why
Moderate load Read from secondaries Simple, effective
Analytics workload Hidden secondary Isolate heavy queries
Massive dataset Sharding Only way to scale beyond single machine
Strong consistency Read from primary No stale reads
Q10 Hard Amazon Meta
How do you scale MongoDB writes? What are the limitations and solutions?

βœ“ Complete Answer:

Write Scaling Challenge:

Unlike reads, ALL writes must go through the primary. This is a fundamental bottleneck in replica sets.

⚠️ Replica Set Write Limitation:

In a replica set, you CANNOT scale writes by adding more secondaries. Secondaries only replicate writes from primary; they don't accept client writes.

Write Scaling Strategies:
1. Vertical Scaling (Short-term)

Upgrade primary to more powerful hardware

  • More CPU cores for parallel processing
  • More RAM for larger working set
  • Faster SSDs for disk I/O
  • Limitation: Eventually hit hardware ceiling
2. Optimize Write Operations
// Batch inserts (much faster than individual)
db.logs.insertMany([...1000 documents...])  // 1 network round trip

// vs individual inserts
for (doc of docs) {
  db.logs.insertOne(doc)  // 1000 network round trips!
}

// Use bulk operations
var bulk = db.logs.initializeUnorderedBulkOp()
bulk.insert({ ... })
bulk.insert({ ... })
bulk.execute()  // Single batch

// Reduce index overhead
// Each index adds write cost
// Keep only necessary indexes
3. Lower Write Concern (Trade-off)
// w:1 - Acknowledge from primary only (fastest)
db.logs.insertOne(
  { event: "user_click" },
  { writeConcern: { w: 1 } }
)

// w:"majority" - Wait for majority (slower but safer)
db.orders.insertOne(
  { total: 1000 },
  { writeConcern: { w: "majority" } }
)

// Trade-off: speed vs durability
4. Sharding - The Real Solution

The ONLY way to truly scale writes horizontally

// Without sharding: 1 primary = 10K writes/sec max

// With 4 shards: 4 primaries = 40K writes/sec
Shard 1 Primary: 10K writes/sec (userId 1-250K)
Shard 2 Primary: 10K writes/sec (userId 250K-500K)
Shard 3 Primary: 10K writes/sec (userId 500K-750K)
Shard 4 Primary: 10K writes/sec (userId 750K-1M)

Each shard handles subset of writes based on shard key
Write Distribution in Sharded Cluster:

Application β†’ mongos (routes based on shard key)
                β”‚
                β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                β–Ό          β–Ό          β–Ό          β–Ό
            Shard 1    Shard 2    Shard 3    Shard 4
            Primary    Primary    Primary    Primary
            
Each primary accepts writes for its data range
Write capacity scales linearly with shard count
5. Application-Level Partitioning

Multiple independent MongoDB clusters

// Route writes to different clusters by geography
if (user.region === "US") {
  writeToCluster("us-mongo-cluster")
} else if (user.region === "EU") {
  writeToCluster("eu-mongo-cluster")
}

// Or by tenant
writeToCluster(`tenant-${tenantId}-cluster`)

When to Shard:

  • Working set no longer fits in RAM
  • Write throughput exceeds single server capacity
  • Dataset size approaching 2-3 TB per replica set
  • Need geographic distribution
  • Queries taking too long even with good indexes
⚠️ Sharding Overhead:

Sharding adds complexity. Don't shard until you need to:

  • Operational overhead (mongos, config servers, balancer)
  • Query complexity increases
  • Cannot change shard key later
  • More moving parts = more failure modes
Q11 Medium Netflix
What is the MongoDB balancer? How does it work and when does it run?

βœ“ Complete Answer:

What is the Balancer?

The balancer is a background process that automatically distributes chunks evenly across shards to maintain balanced data distribution.

Before Balancing (Uneven):
Shard 1: [50 chunks] β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ
Shard 2: [30 chunks] β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ
Shard 3: [20 chunks] β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ

After Balancing (Even):
Shard 1: [33 chunks] β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ
Shard 2: [33 chunks] β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ
Shard 3: [34 chunks] β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ

Balancer moves chunks from overloaded shards to underloaded shards
How the Balancer Works:
  1. Balancer runs on primary config server
  2. Checks chunk distribution every 60 seconds
  3. If imbalance detected (threshold: 8+ chunk difference)
  4. Selects chunks to migrate
  5. Coordinates chunk migration between shards
  6. Updates metadata in config servers
Chunk Migration Process:
Step 1: Balancer identifies imbalance
   Shard A: 50 chunks
   Shard B: 30 chunks
   Difference: 20 > threshold (8)

Step 2: Select chunk to migrate from Shard A

Step 3: Copy chunk data to Shard B
   - Data copied in background
   - Both shards can serve reads/writes during copy

Step 4: Brief lock for final sync
   - Transfer any writes that occurred during copy
   - Takes milliseconds

Step 5: Update metadata
   - Config servers update: chunk now on Shard B
   - mongos starts routing queries to Shard B

Step 6: Delete old data from Shard A
Balancer Configuration:
// Check balancer status
sh.getBalancerState()

// Enable/disable balancer
sh.startBalancer()
sh.stopBalancer()

// Schedule balancer window (e.g., only run at night)
db.settings.update(
  { _id: "balancer" },
  {
    $set: {
      activeWindow: {
        start: "23:00",  // 11 PM
        stop: "06:00"    // 6 AM
      }
    }
  },
  { upsert: true }
)

// Disable balancing for specific collection during migration
sh.disableBalancing("mydb.users")
sh.enableBalancing("mydb.users")
⚠️ Balancer Impact:
  • Resource usage: Chunk migrations consume I/O and network
  • Performance impact: Can slow queries during migration
  • Best practice: Schedule balancing during off-peak hours
  • Monitoring: Watch for excessive migrations (may indicate poor shard key)
πŸ’‘ Interview Tip:

Mention that while the balancer is automatic, production deployments often schedule it for specific time windows to avoid impacting peak traffic.

πŸ’₯ Failure Scenarios & Recovery

Questions about handling failures and disaster recovery.

Q12 Hard Uber LinkedIn
Walk me through a complete datacenter failure scenario. How does MongoDB handle it?

βœ“ Complete Answer:

Scenario Setup:

5-node replica set distributed across 3 datacenters:

Initial State:
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚   DC1 (Primary)    β”‚  β”‚   DC2 (Secondary)  β”‚  β”‚   DC3 (Secondary)  β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚
β”‚   β”‚ Primary  β”‚     β”‚  β”‚   β”‚Secondary1β”‚     β”‚  β”‚   β”‚Secondary2β”‚     β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚                    β”‚
β”‚   β”‚Secondary3β”‚     β”‚  β”‚   β”‚Secondary4β”‚     β”‚  β”‚                    β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚                    β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
     2 members              2 members              1 member

Then: DC1 FAILS completely (power outage, network failure, etc.)
Timeline of Events:
t=0: DC1 Fails
  • Primary and Secondary3 become unreachable
  • Remaining nodes: Secondary1, Secondary2, Secondary4
  • 3 nodes available out of 5 = majority βœ“
t=10s: Failure Detected
  • Heartbeat timeout (default 10 seconds)
  • Surviving secondaries detect primary is down
  • Election process initiated
t=12s: New Primary Elected
  • Secondary1 or Secondary4 becomes primary (higher priority wins)
  • Let's say Secondary1 in DC2 becomes new primary
  • Total downtime: ~12 seconds
After Failover:
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚   DC1 (OFFLINE)    β”‚  β”‚   DC2 (NEW Primary)β”‚  β”‚   DC3 (Secondary)  β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚
β”‚   β”‚  OFFLINE β”‚     β”‚  β”‚   β”‚ PRIMARY  │◄────┼──┼──│Secondary2β”‚     β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚  β”‚                    β”‚
β”‚   β”‚  OFFLINE β”‚     β”‚  β”‚   β”‚Secondary4β”‚     β”‚  β”‚                    β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚  β”‚                    β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                             3 members = majority βœ“
t=12s+: System Operational
  • βœ“ Writes accepted by new primary in DC2
  • βœ“ Reads served by all 3 surviving members
  • βœ“ Replication continues DC2 β†’ DC3
  • βœ— Lost 2 nodes but still have majority
  • ⚠️ Now vulnerable (losing 1 more node = no majority)
When DC1 Recovers:
1. Nodes in DC1 come back online
2. They discover they're behind new primary
3. Automatic sync from current primary
4. May involve rollback if they had un-replicated writes
5. Rejoin replica set as secondaries
6. System back to full 5-node redundancy
⚠️ Critical: Majority Requirement

Why 3 datacenters?

  • 2 DCs with 2-3 nodes each: Either DC failure = no majority
  • 3 DCs with distributed nodes: Any single DC failure = majority survives
Best Practices for Multi-DC Deployment:
  • Odd number of nodes: 3, 5, or 7 members
  • Majority in primary DC: Keeps primary local during normal operation
  • Priority configuration: Prefer certain DCs for primary
  • Monitoring: Alert immediately on datacenter failures
  • Runbooks: Clear procedures for DC recovery
Q13 Hard Google Amazon
What is a "split-brain" scenario? How does MongoDB prevent it?

βœ“ Complete Answer:

What is Split-Brain?

Split-brain occurs when a network partition causes multiple nodes to believe they are the primary, leading to divergent data and conflicts.

Split-Brain Scenario (DANGEROUS):

Network Partition:
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚    Partition A         β”‚     β”‚     β”‚    Partition B         β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚     β”‚     β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚
β”‚   β”‚PRIMARY A │◄─writes β”‚     β”‚     β”‚   β”‚PRIMARY B │◄─writes β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚     β”‚     β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚     β”‚     β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚
β”‚   β”‚Secondary1β”‚         β”‚     β”‚     β”‚   β”‚Secondary2β”‚         β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚     β”‚     β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚     β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                               β”‚
                        Network Partition

Problem: TWO primaries accepting writes!
         Data diverges β†’ consistency violated!
         Which data is correct when network heals?

MongoDB's Prevention: Majority Voting

How MongoDB Prevents Split-Brain:
  1. Majority Requirement: Primary must see majority of nodes
  2. Heartbeat Monitoring: Primary checks it can reach majority
  3. Step Down: If primary loses majority, it IMMEDIATELY steps down
  4. No Writes: Stepped-down node becomes read-only
  5. Single Primary: Only partition with majority can elect primary
MongoDB's Actual Behavior (SAFE):

Network Partition:
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”‚     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚    Partition A (2/5)   β”‚     β”‚     β”‚    Partition B (3/5)   β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚     β”‚     β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚
β”‚   β”‚SECONDARY β”‚         β”‚     β”‚     β”‚   β”‚ PRIMARY  │◄─writes β”‚
β”‚   β”‚(stepped) β”‚         β”‚     β”‚     β”‚   β”‚(elected) β”‚         β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚     β”‚     β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚
β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚     β”‚     β”‚   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”‚
β”‚   β”‚Secondary1β”‚         β”‚     β”‚     β”‚   β”‚Secondary2β”‚         β”‚
β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚     β”‚     β”‚   β”‚Secondary3β”‚         β”‚
β”‚                        β”‚     β”‚     β”‚   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β”‚
β”‚   NO MAJORITY          β”‚     β”‚     β”‚   HAS MAJORITY βœ“       β”‚
β”‚   Read-only mode       β”‚     β”‚     β”‚   Read + Write         β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β”‚     β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Result: Only ONE primary, in partition with majority
        No split-brain possible!
Step-by-Step Prevention:
5-node replica set: Primary + 4 Secondaries

Network partition splits them:
- Partition A: 2 nodes
- Partition B: 3 nodes

Partition A (has old primary):
1. Primary detects it can only reach 1 other node (2 total)
2. 2 < 3 (majority of 5) β†’ NO MAJORITY
3. Primary STEPS DOWN immediately
4. Becomes SECONDARY (read-only)
5. Cannot elect new primary (need 3 votes, only have 2)

Partition B:
1. Detects primary is gone
2. Has 3 nodes = MAJORITY βœ“
3. Initiates election
4. Elects NEW primary
5. Accepts writes

Result: Only partition B has writable primary
        Partition A is read-only
        No conflicting writes = No split-brain!
Why This Works:

Mathematical Guarantee:

  • In any partition, at MOST one side can have majority
  • 5 nodes split any way: max one side has β‰₯3 nodes
  • Example splits:
    • 1-4 split: 4 > 2.5 βœ“ (only 4-side has majority)
    • 2-3 split: 3 > 2.5 βœ“ (only 3-side has majority)
    • 3-2 split: 3 > 2.5 βœ“ (only 3-side has majority)
  • No way to get 2 sides both with majority!
⚠️ Even Number Anti-Pattern:
4-node replica set (BAD):

2-2 network partition:
- Partition A: 2 nodes (not majority: 2 < 2.5)
- Partition B: 2 nodes (not majority: 2 < 2.5)

Result: NEITHER side can elect primary!
        Entire system becomes read-only!
        
Solution: Use 5 nodes (odd number) or add arbiter to make 5 voting members
πŸ’‘ Interview Gold:

Explain that MongoDB's majority requirement isn't just good designβ€”it's a mathematical guarantee that split-brain is IMPOSSIBLE. This shows deep understanding of distributed systems theory.

πŸ—οΈ System Design Questions

Real-world architecture and design decisions.

Q14 Hard Meta Uber
Design a MongoDB architecture for a social media platform with 100M users. Consider read/write patterns, data distribution, and failure scenarios.

βœ“ Complete Answer:

Requirements Analysis:

Scale Characteristics:
  • 100M users, ~10M DAU (daily active)
  • Heavy read workload (90% reads, 10% writes)
  • User-generated content: posts, comments, likes
  • Social graph: followers/following
  • Feed generation (personalized, real-time)
  • Global user base (latency-sensitive)

Proposed Architecture:

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚                     Application Layer                        β”‚
β”‚                   (Load Balanced Servers)                    β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                           β”‚
                           β–Ό
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚                   Redis Cache Layer                          β”‚
β”‚         (Hot data, session, feed cache)                      β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                           β”‚
                           β–Ό
β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚                   MongoDB Sharded Cluster                    β”‚
β”‚                                                              β”‚
β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   β”‚
β”‚  β”‚ Shard 1  β”‚  β”‚ Shard 2  β”‚  β”‚ Shard 3  β”‚  β”‚ Shard 4  β”‚   β”‚
β”‚  β”‚Users     β”‚  β”‚Users     β”‚  β”‚Users     β”‚  β”‚Users     β”‚   β”‚
β”‚  β”‚1-25M     β”‚  β”‚25-50M    β”‚  β”‚50-75M    β”‚  β”‚75-100M   β”‚   β”‚
β”‚  β”‚          β”‚  β”‚          β”‚  β”‚          β”‚  β”‚          β”‚   β”‚
β”‚  β”‚(P+2S)    β”‚  β”‚(P+2S)    β”‚  β”‚(P+2S)    β”‚  β”‚(P+2S)    β”‚   β”‚
β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜   β”‚
β”‚                                                              β”‚
β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”                 β”‚
β”‚  β”‚ Shard 5  β”‚  β”‚ Shard 6  β”‚  β”‚ Shard 7  β”‚                 β”‚
β”‚  β”‚Posts     β”‚  β”‚Posts     β”‚  β”‚Posts     β”‚                 β”‚
β”‚  β”‚          β”‚  β”‚          β”‚  β”‚          β”‚                 β”‚
β”‚  β”‚(P+4S)    β”‚  β”‚(P+4S)    β”‚  β”‚(P+4S)    β”‚                 β”‚
β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                 β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
Collection Design:
// Users Collection (Shard Key: userId)
{
  _id: ObjectId,
  userId: Long,  // Shard key
  username: String,
  email: String,
  profile: {
    name: String,
    bio: String,
    avatar: String
  },
  stats: {
    followers: Number,
    following: Number,
    posts: Number
  },
  createdAt: Date
}

// Posts Collection (Shard Key: userId + timestamp)
{
  _id: ObjectId,
  postId: Long,
  userId: Long,  // Part of shard key
  content: String,
  mediaUrls: [String],
  likes: Number,
  comments: Number,
  timestamp: Date,  // Part of shard key
  tags: [String]
}

// Social Graph Collection (Shard Key: userId)
{
  _id: ObjectId,
  userId: Long,  // Shard key
  followers: [Long],     // Array of userIds
  following: [Long],     // Array of userIds
  lastUpdated: Date
}
Sharding Strategy:
Users Collection:
sh.shardCollection("social.users", { userId: "hashed" })

Why hashed?
βœ“ Even distribution of users across shards
βœ“ Prevents hotspots from sequential user IDs
βœ“ Queries by userId are targeted (single shard)

Replica Set Configuration:
- 3 nodes per shard (1 Primary + 2 Secondaries)
- Geographic distribution across 3 datacenters
- Read preference: primaryPreferred for user profiles
Posts Collection:
sh.shardCollection("social.posts", { userId: 1, timestamp: 1 })

Why compound key?
βœ“ userId ensures user's posts stay together (feed queries)
βœ“ timestamp prevents hotspots on recent posts
βœ“ Range queries on timestamp within user work well

Replica Set Configuration:
- 5 nodes per shard (1 Primary + 4 Secondaries)
- More secondaries for high read load (feed generation)
- Read preference: secondary for feed reads
- Write concern: w:1 for posts (eventually consistent OK)
Read Optimization:
  • Caching Layer: Redis for hot user data, feed cache (TTL 5 min)
  • Read from Secondaries: Offload feed generation to secondaries
  • Materialized Views: Pre-compute popular user feeds
  • CDN: Media files served from CDN, not MongoDB
Write Optimization:
  • Write Concern w:1: Fast acknowledgment for posts/likes
  • Bulk Operations: Batch notifications, likes
  • Async Processing: Feed updates via message queue
High Availability:
Geographic Distribution:
- Primary DC: 40% of nodes
- Secondary DC 1: 30% of nodes  
- Secondary DC 2: 30% of nodes

Failure Handling:
βœ“ Single shard failure: Automatic failover (12s downtime)
βœ“ Datacenter failure: Majority survives, operations continue
βœ“ Config server: 3-node replica set for metadata
βœ“ Mongos: Stateless routers, run multiple instances
Monitoring & Alerts:
  • Replication lag > 1 second
  • Chunk migration activity
  • Shard imbalance > 10 chunks
  • Write hotspots (single shard > 50% writes)
  • Connection pool exhaustion
Expected Performance:
  • Reads: 100K+ QPS (distributed across secondaries)
  • Writes: 10K+ QPS (distributed across shard primaries)
  • Latency: p99 < 100ms for reads, < 50ms for writes
  • Availability: 99.95% uptime (4 hours downtime/year)
🎯 Interview Tip:

Walk through trade-offs: Why sharding over replication? Why hashed vs range shard keys? Why different replica set sizes for different collections? Show you can make decisions based on access patterns.