π 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.
- 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
- 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.
β 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
]
})
π Follow-up Topics:
- How does the election process work when primary fails?
- What is the oplog and how does it work?
- How do you configure read preferences?
β 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)
- 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)
]
}
π Follow-up Questions:
- What happens to in-flight writes during election?
- How do you prevent elections during maintenance?
- What is "rollback" and when does it occur?
β 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 |
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.
βοΈ Consistency & CAP Theorem
Questions about data consistency and distributed systems trade-offs.
β 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"
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.
β 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
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.
β 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:
- Old primary detects it's behind new primary
- Finds common point in oplog
- Saves rolled-back operations to rollback files
- Reverts to common point
- Syncs from new primary
- 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
- 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
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.
β 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 |
β Complete Answer:
Write Scaling Challenge:
Unlike reads, ALL writes must go through the primary. This is a fundamental bottleneck in replica sets.
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 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
β 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:
- Balancer runs on primary config server
- Checks chunk distribution every 60 seconds
- If imbalance detected (threshold: 8+ chunk difference)
- Selects chunks to migrate
- Coordinates chunk migration between shards
- 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")
- 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)
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.
β 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
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
β 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:
- Majority Requirement: Primary must see majority of nodes
- Heartbeat Monitoring: Primary checks it can reach majority
- Step Down: If primary loses majority, it IMMEDIATELY steps down
- No Writes: Stepped-down node becomes read-only
- 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!
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
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.
β 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)
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.