Type something to search...
MongoDB Oplog Explained: How Replication Works Under the Hood

MongoDB Oplog Explained: How Replication Works Under the Hood

Most of the time, replication in MongoDB is invisible. You write to the primary, the secondaries catch up, and everything just works. Then one day a secondary falls into RECOVERING and refuses to come back without a full resync, or a nightly batch job causes replication lag to spike to twenty minutes, or someone asks why a restore window is only six hours long. All of those questions have the same answer: the oplog.

The oplog (operations log) is a special collection that records every change made to the data on a replica set, in order. Secondaries replicate by copying and replaying it. Change streams read from it. Point-in-time recovery depends on it. Once you understand what goes into the oplog, how big it is, and how secondaries consume it, a whole category of production mysteries becomes easy to reason about.

This guide covers where the oplog lives, what an entry looks like, why operations are rewritten to be idempotent, how secondaries sync and apply entries, how to size and monitor the oplog window, and what happens during initial sync and rollback.

Where the Oplog Lives

Every member of a replica set has an oplog: the collection oplog.rs in the local database. The local database is special, because it isn't replicated. Each member maintains its own copy, which is what makes it safe to use as the replication log itself.

You can look at it directly from mongosh:

use local
db.oplog.rs.stats().maxSize   // configured maximum size in bytes
db.oplog.rs.find().sort({ $natural: -1 }).limit(1)   // the most recent entry

The oplog behaves like a capped collection: it has a fixed maximum size, entries are stored in insertion order, and when it's full the oldest entries are removed to make room for new ones. Two refinements make it safer than a plain capped collection. It can temporarily grow beyond its configured size so that it never deletes entries a majority of members still need, and you can set a minimum retention period (covered below) so entries are kept for a guaranteed amount of time regardless of size.

Anatomy of an Oplog Entry

Insert a document and look at the resulting entry:

use shop
db.orders.insertOne({ _id: 1001, status: "paid", total: 84.5 })

use local
db.oplog.rs.find({ ns: "shop.orders" }).sort({ $natural: -1 }).limit(1)
{
  op: "i",
  ns: "shop.orders",
  ui: UUID("3f5c2a8e-7b41-4d6a-9f0e-1c2b3a4d5e6f"),
  o: { _id: 1001, status: "paid", total: 84.5 },
  o2: { _id: 1001 },
  ts: Timestamp({ t: 1726908060, i: 1 }),
  t: Long(7),
  v: Long(2),
  wall: ISODate("2026-09-21T08:41:00.214Z")
}

The fields you'll care about:

FieldMeaning
opOperation type: i insert, u update, d delete, c command, n no-op
nsNamespace (database.collection)
uiThe collection's UUID, stable even if the collection is renamed
oThe operation itself: the document, the update diff, or the command
o2For updates, the _id of the target document
tsTimestamp: seconds plus an increment, unique and ordered across the set
tElection term in which the primary wrote this entry
wallWall-clock time of the write

The ts field is the backbone of replication. It's a BSON Timestamp, a 64-bit value made of a Unix time in seconds and an ordinal counter for operations within that second. Every entry's ts is strictly greater than the one before it, so members can talk about "I've applied everything up to ts X" with no ambiguity. Combined with the term t, it forms an optime, which is what you see in rs.status().

Operations Are Rewritten to Be Idempotent

Here's the most important design rule of the oplog: applying an entry once or many times must produce the same result. Secondaries may re-apply entries during recovery or initial sync, and replaying $inc: { stock: -1 } twice would corrupt data.

So the primary doesn't log your update as you wrote it. It logs the effect of the update. Watch what happens to an $inc:

use shop
db.orders.updateOne({ _id: 1001 }, { $inc: { total: 10 }, $set: { status: "shipped" } })

use local
db.oplog.rs.find({ ns: "shop.orders", op: "u" }).sort({ $natural: -1 }).limit(1)
{
  op: "u",
  ns: "shop.orders",
  o: { $v: 2, diff: { u: { total: 94.5, status: "shipped" } } },
  o2: { _id: 1001 },
  ts: Timestamp({ t: 1726908121, i: 1 }),
  // ...
}

The $inc became a plain "set total to 94.5." Replaying that entry any number of times leaves total at 94.5. In recent versions, updates are logged in this compact diff format ($v: 2), where u lists updated fields, i lists inserted fields, d lists deleted fields, and keys prefixed with s describe changes inside nested subdocuments.

The same rule affects multi-document operations. A single deleteMany that removes 50,000 documents doesn't produce one oplog entry; it produces 50,000 delete entries, one per _id:

db.oplog.rs.find({ ns: "shop.sessions", op: "d" }).limit(2);
{ op: "d", ns: "shop.sessions", o: { _id: ObjectId("66ee...a1") }, ts: Timestamp({ t: 1726908200, i: 1 }), ... }
{ op: "d", ns: "shop.sessions", o: { _id: ObjectId("66ee...a2") }, ts: Timestamp({ t: 1726908200, i: 2 }), ... }

This is why bulk updates and deletes are so much more expensive for replication than they look. One command on the primary can generate millions of oplog entries, all of which every secondary must fetch and apply, and all of which consume oplog space.

Commands, Transactions, and No-ops

Not everything in the oplog is a document change:

  • op: "c" records commands like create, drop, createIndexes, and renameCollection. Multi-document transactions also appear as command entries with an applyOps array containing the individual operations, so the whole transaction is applied atomically on secondaries.
  • op: "n" entries are no-ops. The primary writes one periodically when it's otherwise idle, so the latest optime keeps advancing. That lets secondaries and change streams know the clock is still moving even when nobody is writing.

How Secondaries Replicate

Each secondary runs a replication loop with three moving parts:

  1. Fetching. The secondary opens a tailable cursor on its sync source's oplog, asking for everything after the last ts it has. The sync source is usually the primary, but it can be another secondary, a setup called chained replication, which reduces load on the primary and is useful across data centers.
  2. Buffering. Fetched entries go into an in-memory buffer.
  3. Applying. Batches of entries are applied to the secondary's data and written into its own oplog. Application is multi-threaded: entries for different documents can be applied in parallel, while entries for the same document are applied in order.

The secondary's position is reported back to the primary, which is how the primary knows when a write has reached a majority. That feeds directly into write concern: a write with w: "majority" is acknowledged when enough members report an optime at or beyond it.

You can see each member's sync source and position in rs.status():

rs.status().members.map((m) => ({
  name: m.name,
  state: m.stateStr,
  syncSource: m.syncSourceHost,
  optime: m.optime?.ts,
}));
[
  {
    name: "db-1:27017",
    state: "PRIMARY",
    syncSource: "",
    optime: Timestamp({ t: 1726908300, i: 4 }),
  },
  {
    name: "db-2:27017",
    state: "SECONDARY",
    syncSource: "db-1:27017",
    optime: Timestamp({ t: 1726908300, i: 4 }),
  },
  {
    name: "db-3:27017",
    state: "SECONDARY",
    syncSource: "db-2:27017",
    optime: Timestamp({ t: 1726908299, i: 1 }),
  },
];

Here db-3 is syncing from db-2 (chained) and is about one second behind.

Where Replication Lag Comes From

Lag is the gap between the primary's latest optime and a secondary's. The common causes map directly to the loop above:

  • Fetching is slow: network bandwidth or latency between the secondary and its sync source.
  • Applying is slow: the secondary has slower disks, a smaller cache, or is busy serving heavy reads.
  • The primary produces entries faster than they can be applied: usually a bulk job generating huge numbers of per-document entries.
  • Unindexed updates on secondaries: oplog entries target documents by _id, so this is less of an issue than it used to be, but index builds and large collection scans on the secondary still compete for resources.

rs.printSecondaryReplicationInfo() gives a quick per-member lag readout, and your monitoring should track it continuously.

Sizing the Oplog and the Oplog Window

The oplog window is the time span between the oldest and newest entries in the oplog. It's the most important number to know about your oplog, because it defines how long a secondary can be offline and still catch up by replaying entries.

rs.printReplicationInfo();
actual oplog size
'10240 MB'
---
configured oplog size
'10240 MB'
---
log length start to end
'154812 secs (43 hrs)'
---
oplog first event time
'Sat Sep 19 2026 13:40:12 GMT+0000'
---
oplog last event time
'Mon Sep 21 2026 08:40:24 GMT+0000'
---
now
'Mon Sep 21 2026 08:41:02 GMT+0000'

A 43-hour window means a member can be down for up to 43 hours and rejoin without a full resync. The window isn't fixed: it depends on write volume. If a batch job doubles your write rate for a night, the window shrinks accordingly.

By default, MongoDB sizes the oplog at 5% of free disk space, with a floor of 990 MB and a cap of 50 GB, on WiredTiger. For write-heavy workloads, that default is often too small. You can resize it online, on each member, without a restart:

// Size is in megabytes; run on each member
db.adminCommand({ replSetResizeOplog: 1, size: 20480 });

// Guarantee entries are kept for at least 48 hours, even if that exceeds the size
db.adminCommand({ replSetResizeOplog: 1, minRetentionHours: 48 });

The minimum retention option is valuable: it turns the oplog window from "whatever the current write rate allows" into a guaranteed floor, at the cost of the oplog growing larger during write bursts. Make sure you have the disk headroom before enabling it.

A good target: an oplog window of at least 24 to 72 hours, comfortably longer than your longest maintenance, your slowest initial sync, and the maximum time a change stream consumer might be down.

Initial Sync and Full Resyncs

When a brand-new member joins, or a member falls so far behind that the entries it needs have been deleted from every sync source's oplog, it performs an initial sync:

  1. It records the sync source's current oplog position.
  2. It clones every database and collection (except local) from the sync source and builds indexes.
  3. It applies all oplog entries written during the clone, starting from the recorded position.
  4. It transitions to SECONDARY once it's consistent.

Step 3 is why the oplog window must be longer than an initial sync takes. If cloning a 2 TB dataset takes 20 hours and the oplog window is 12 hours, the entries needed in step 3 are gone before the clone finishes, and the sync fails and starts over. On large datasets, seeding a new member from a filesystem snapshot or backup is often faster than a logical initial sync.

A member that has fallen off the end of the oplog reports something like this in its log and sits in RECOVERING:

"msg":"We are too stale to use as a sync source...","attr":{"lastOpTimeFetched":...,"earliestOpTimeAvailable":...}

At that point the fix is a resync, which is exactly the situation a well-sized oplog prevents.

Rollbacks: When the Oplog Diverges

Suppose the primary accepts a write with w: 1, then crashes before any secondary fetches it. A secondary is elected, and new writes continue on it. When the old primary comes back, its oplog contains an entry that the rest of the set never saw.

The old primary finds the last optime it has in common with the new primary, undoes everything after that point, and writes the undone documents to rollback files in its data directory (under rollback/), so an operator can inspect and manually re-apply them if needed. Then it resumes normal replication.

Rollbacks only affect writes that weren't replicated to a majority. That's the practical case for w: "majority": a write acknowledged by a majority is guaranteed to be on whichever member wins the next election, so it can never be rolled back. The deeper trade-offs between acknowledgement levels are covered in MongoDB Read Preferences and Write Concerns Explained.

The Oplog and Change Streams

Before change streams existed, applications tailed local.oplog.rs directly to react to data changes. That approach has serious problems: the entry format is internal and has changed between versions (the update diff format is one example), it exposes writes that haven't been majority-committed and might roll back, and it doesn't work cleanly across shards.

Change streams are the supported interface. They're built on the oplog, but they only deliver majority-committed changes, present a stable documented event format, work across sharded clusters, and give you resumable tokens. If you're tempted to read the oplog from application code, use a change stream instead. For a worked example, see Building an Event-Driven Audit Log with MongoDB Change Streams.

Change streams do inherit the oplog's limits, though. A consumer that's offline longer than the oplog window can't resume, which is one more reason to size the window generously.

Common Pitfalls

Running huge multi-document writes during peak hours. A single updateMany over 20 million documents generates 20 million oplog entries, drains the oplog window, and spikes replication lag. Batch large changes and throttle them, watching lag as you go.

Assuming the default oplog size is enough. The default was chosen for typical workloads, not yours. Check the window with rs.printReplicationInfo() during your busiest period, not on a quiet Sunday.

Forgetting that TTL deletes and bulk imports count too. A TTL index expiring millions of documents a day, or a nightly import, can dominate oplog traffic. Include them when you estimate the window.

Writing to the local database. Never insert, update, or delete in local.oplog.rs. Manual changes can corrupt replication on that member.

Ignoring rollback directories. If a rollback happens, the data it removed sits in rollback files until someone looks. Alert on rollback events and review the files.

Letting secondaries run on weaker hardware. Secondaries do the same write work as the primary. If they have slower disks, they'll lag under load and take longer to become primary during a failover.

Conclusion

The oplog is an ordered, idempotent record of every change on a replica set, stored in local.oplog.rs on each member. Secondaries fetch it from a sync source, apply it in parallel, and report their progress so the primary can honor majority write concern. Its size defines the oplog window, which in turn decides whether a lagging member or change stream consumer can catch up or needs a full resync, and majority writes keep the oplog from ever diverging in ways that cause rollbacks.

Next step: run rs.printReplicationInfo() on your production primary during peak traffic and note the window. If it's under 24 hours, resize the oplog with replSetResizeOplog or set minRetentionHours, and add an alert for when the window drops below your threshold.

Tags :
Share :

Related Posts

A Complete Guide to MongoDB Query Operators

A Complete Guide to MongoDB Query Operators

Your first MongoDB queries are usually simple equality filters: find the user with this email, find orders with this status. That covers a surprising

Continue Reading
Async MongoDB in Python with Motor and FastAPI

Async MongoDB in Python with Motor and FastAPI

FastAPI runs your endpoints on an event loop. That's what lets a single worker juggle hundreds of concurrent requests: while one request waits on the

Continue Reading
Atlas Online Archive: Tiering Cold Data to Cut Costs

Atlas Online Archive: Tiering Cold Data to Cut Costs

Look at almost any production database and you'll find the same shape. A small slice of recent data gets nearly all the reads and writes: this week's

Continue Reading