Lesson 0007 · Replication

Replica sets

Mongo's high-availability primitive — automatic failover, not a bolt-on. One primary, N secondaries, one oplog. The oplog is what makes it tick: a capped, logical, idempotent operation log that your Postgres and MySQL courses will recognize immediately.

~12 minWarm-up + retrieval quiz + replication contrastCheat sheet: here

Warm-up — apply lesson 0006 first (cold, from memory)

One device, potentially millions of sensor readings. Where does the parent reference go, and why?

One-to-squillions: an array of reading _ids on the device document grows forever and eventually hits the 16 MB limit. Flip the reference — each reading carries a deviceId. The device document stays small; querying all readings for a device is an indexed range scan on the readings collection. Junction collections are a SQL pattern; MongoDB doesn't need them.

You embed an author sub-document inside the blog post document. Updating both the post title and the author's bio in one operation is atomic because:

Single-document atomicity: any write to one document — including all nested fields and arrays — is all-or-nothing. If the author were in a separate referenced document you'd need two writes and either inconsistency or an explicit multi-document transaction.

Why is embedding event-log entries in an array on the parent document an anti-pattern?

The unbounded array anti-pattern: every new event extends the array. The document grows monotonically, eventually hitting the 16 MB cap on a write. Well before that limit, every read of the parent pulls in stale log data you didn't ask for. Reference on the child side instead.

In Postgres you configure streaming replication manually and promote standbys by hand (or outsource it to Patroni). In MongoDB, high availability is built into the topology primitive: a replica set. There is no "standalone plus optional replication" mode for production — a replica set is the default deployment, always.

Architecture: one primary, N secondaries

A replica set is a group of mongod instances — typically three — that hold the same data set. — MongoDB Manual: "A replica set in MongoDB is a group of mongod processes that maintain the same data set. Replica sets provide redundancy and high availability." Exactly one member is the primary at any moment; the rest are secondaries. All writes go to the primary; secondaries replicate by tailing the primary's operation log (the oplog).

CLIENT │ (all writes, default reads) ▼ PRIMARY ──── oplog ────► SECONDARY 1 └──► SECONDARY 2

A third node type is the arbiter: a lightweight process with no data, which exists only to vote in elections. Arbiters let you run a 2-data-node set (cheaper) with a quorum still reachable. They add no redundancy — don't use them if you can afford three full nodes. — MongoDB Manual: "An arbiter does not have a copy of data set and cannot become a primary… Add an arbiter to a replica set to have an odd number of members."

The oplog — the heart of replication

The oplog lives at local.oplog.rs — a capped collection in the local database on every replica set member. — MongoDB Manual: "The oplog (operations log) is a special capped collection that keeps a rolling record of all operations that modify the data stored in your databases." Every write the primary executes is recorded here as a logical, idempotent operation. Secondaries tail this collection and apply each entry in order.

Two properties matter most:

Each oplog entry is a BSON document. The key fields:

ts — timestamp (used for ordering and for replica lag calculation) op — operation type: "i" (insert), "u" (update), "d" (delete), "c" (command) ns — namespace: "dbName.collectionName" o — the operation body (the full document for inserts; the update spec for updates)

The contrast with your existing courses:

EngineReplication logLevelIdempotent?
Postgres streamingWAL (write-ahead log)Physical — byte-level page changesNo (page apply is not idempotent)
MySQL binlog (row-based)Binary logLogical — before/after row imagesIn practice yes (with row images)
MongoDB replica setoplog (local.oplog.rs)Logical — full transformationYes, by design

The physical/logical distinction matters for portability and recovery. Postgres WAL is tied to the exact on-disk page format and the major version — you can't stream a Postgres 14 WAL to a Postgres 15 standby. MongoDB's logical oplog is engine-independent: the same oplog entry can be applied regardless of what WiredTiger internal pages look like. — MongoDB Manual: "MongoDB applies database operations on the primary and then records the operations on the primary's oplog. The secondary members then copy and apply these operations in an asynchronous process."

Elections — the built-in failover

Members send each other heartbeats every two seconds. If a secondary doesn't hear from the primary within electionTimeoutMillis (default 10 seconds), it concludes the primary is unreachable and calls for an election. — MongoDB Manual: "When the primary is unavailable, an eligible secondary will hold an election to select itself as the new primary. The first secondary to call the election and receive votes from a majority of the voting members becomes the new primary."

The winner must:

  1. Have the most up-to-date oplog position (the highest ts) among the candidates.
  2. Receive votes from a majority of voting members — this is where the 3-node recommendation comes from. A 3-node set (majority = 2) can survive one failure. A 2-node set without an arbiter cannot elect a new primary (can't reach majority of 2).

During an election — typically 10–30 seconds — writes fail (no primary to accept them) and reads from primary also fail. This is the availability window your SLA must account for. Applications should use drivers with retryable writes enabled (MongoDB driver default since 4.2) so transient election failures are hidden.

Election vocabulary your Postgres course touches: quorum (here it's majority of voting members, not just nodes); priority (you can set priority 0 on a member to prevent it ever becoming primary — useful for a geographically distant data-center replica you want to keep as read replica only).

Reading from secondaries

By default, all reads go to the primary (readPreference: "primary") — you always read your own writes. You can change this per operation or per connection: secondaryPreferred routes to a secondary if one is available, falling back to primary; secondary always routes to a secondary, accepting potential staleness. — MongoDB Manual: "Read preference describes how MongoDB clients route read operations to the members of a replica set… By default, an application directs its read operations to the primary member in a replica set." Stale reads are the price: a secondary may be seconds behind the primary's oplog. Lesson 0008 (write/read concern) covers how to trade consistency and durability precisely.

Bridge to 0008: How durable is a write? With the default w: 1, the primary acknowledges as soon as it writes to its own oplog — secondaries may not have applied it yet. With w: "majority", it waits for a majority of voting members to confirm. Combined with j: true (journal flush before ack), you control the full durability/latency trade-off. That's the whole of lesson 0008.

A 3-node replica set can survive how many simultaneous node failures while still electing a new primary?

Elections require a majority of voting members — 2 of 3. If two nodes fail, only 1 of 3 remains; it cannot reach majority and cannot elect a primary. Writes stop. This is why 3 is the minimum recommended size for a fault-tolerant replica set.

The oplog is idempotent. What does that guarantee if a secondary crashes mid-apply and replays the same oplog entry on restart?

Idempotency means applying f(x) twice = applying f(x) once. MongoDB stores operations as full transformations, so re-applying the same oplog entry after a crash produces the same document state. No deduplication mechanism is needed — the engine can just re-run the oplog from the last confirmed position.

The MongoDB oplog vs the Postgres WAL — the sharpest distinction:

Postgres WAL records byte-level changes to on-disk pages — it's version-tied and engine-specific. MongoDB's oplog records logical operations (insert/update/delete as BSON documents) — engine-independent and human-readable. MySQL's row-based binlog is also logical but stores before/after images rather than idempotent transformations.

A secondary's oplog position has fallen behind a point that has been overwritten (the oplog window has passed). What happens next?

The oplog is a capped collection — oldest entries are overwritten when full. A secondary that falls outside the oplog window has lost its place and can't replay the missing ops. It must re-copy the entire data set from a healthy member. The direct Postgres analog: a lagging replica outrunning the WAL segment retention window.

By default, where does a MongoDB client send read operations in a replica set?

The default readPreference is "primary" — all reads go to the primary, giving strong consistency (you always read the latest committed write). Routing reads to secondaries is opt-in (secondaryPreferred, secondary) and accepts the possibility of stale reads because secondaries trail the primary's oplog asynchronously.
Optional — when you have a mongosh instance: start three mongod processes locally and initiate a replica set with rs.initiate(). Run rs.status() to see the primary/secondary assignments and oplog lag. Then db.printReplicationInfo() to see the oplog window (how many hours of operations it can hold). Kill the primary process and watch rs.status() as an election fires.