The idea: from a command to a WiredTiger page, and out to the replica set

A MongoDB command passes through three parts. mongod receives the command, plans it and runs it. The WiredTiger storage engine keeps the data in B-trees, in memory and on disk. The replica set copies every write to other servers. The animation follows one command at a time through all three. It uses one client, a 3-member replica set (A is the primary, B and C are secondaries) and the defaults of MongoDB 8.0.

The driver sends an OP_MSG to mongod, whose transport layer, command dispatch, query planner and plan executor pass the command down to WiredTiger. WiredTiger keeps collection-0, the _id_ and age_1 indexes and the oplog as B-trees in its cache, with a journal and checkpointed .wt files on disk. Primary A sends its oplog to secondaries B and C.
A command goes down through mongod's layers into WiredTiger's cached B-trees, and every write also flows out through the oplog to the secondaries.
db.users.insertMany([
  {_id: 1, name: "Ann", age: 25}, {_id: 2, name: "Bob", age: 32}, {_id: 3, name: "Cid", age: 28},
  {_id: 5, name: "Eve", age: 30}, {_id: 6, name: "Fay", age: 22}, {_id: 7, name: "Gus", age: 30},
  {_id: 9, name: "Ivy", age: 27}, {_id: 10, name: "Jon", age: 35}])
db.users.createIndex({age: 1})      // "age_1"; the "_id_" index always exists

These are the same 8 rows as on the MySQL and PostgreSQL pages, so you can compare the three.

Reading the canvas: in the WiredTiger cache a grey page is only on disk (not in the cache), an orange page is dirty (changed since the last checkpoint), ✗ marks a tombstone and ← an older version on a record's update chain. In the replica set, A has priority 2 and B and C priority 1; with three members, 2 of 3 make a majority.

The layers of mongod

LayerWhat it does
driverTurns db.users.find(…) into a BSON command such as {find: "users", filter: {…}}. Sends it in an OP_MSG message over a pooled connection to the primary.
transport layerOne thread per connection reads the message.
command dispatchFinds the command, checks authorization, and takes intent locks (IS for reads, IX for writes) on the global, database and collection levels. Documents are never locked here.
query plannerChooses the stages that run the query: EXPRESS_IXSCAN, IXSCAN → FETCH or COLLSCAN. Uses the plan cache when it can.
plan executorRuns the stages. Each stage asks WiredTiger cursors for the next key or record.
WiredTigerB-tree tables in a cache, transactions and snapshots (MVCC), the journal and checkpoints.

Documents, RecordIds and the _id index

A collection is stored as a WiredTiger table (collection-0 here). Its key is an internal 64-bit RecordId and its value is the BSON document. RecordIds are handed out in increasing order, so new documents land in the rightmost leaf. The _id field is not the storage key. It is a key in a separate unique index, _id_ (index-1), which maps _id → RecordId. Every other index works the same way. age_1 stores the key (age, RecordId): adding the RecordId keeps equal ages distinct and in order.

So every index lookup is index → RecordId → document. (Collections created with the clusteredIndex option store documents by _id directly, like InnoDB's clustered index. Those are not shown.) If _id is missing, the driver adds an ObjectId: a 4-byte timestamp, 5 random bytes and a 3-byte counter.

Three columns: the _id_ index maps _id 1 to 10 onto RecordIds r1 to r8, collection-0 maps each RecordId to its BSON document, and age_1 holds (age, RecordId) pairs in age order. find({age: 30}) scans age_1 entries (30, r4) and (30, r6) and fetches the documents of Eve and Gus.
Every index, _id included, stores RecordIds; a query scans an index and then fetches each document by RecordId.

The query planner and the plan cache

  • Equality on _id skips planning. MongoDB 8.0 calls it the express path (EXPRESS_IXSCAN); older versions called it IDHACK. (Demo: find by _id)
  • An indexed field gives an IXSCAN over the index bounds. Above it, a FETCH reads each document by RecordId. In explain(), keysExamined and docsExamined count this work. (Demo: find through the age_1 index)
  • No usable index gives a COLLSCAN: every document is read and tested. (Demo: COLLSCAN)
  • When several indexes fit, the planner runs every candidate plan for a short trial. The winner is the plan that returns results with the fewest "works". It is cached in the plan cache under the query's shape: the fields and operators, not the values. The next query of the same shape reuses it.

WiredTiger cache, update chains and MVCC

WiredTiger never overwrites a record in place. Each record in an in-memory page has an update chain of versions, newest first:

  • an update adds a new value (a small WT_MODIFY delta when only a few bytes change);
  • a delete adds a tombstone.

Each version carries its commit timestamp. A reader works in a snapshot and walks the chain to the newest version it may see. Readers never wait for writers. Old versions are dropped once no snapshot can need them. Before that, if memory is short, they are moved to the history store. When a page is written to disk (reconciliation, at a checkpoint or on eviction), only its newest committed state is written. So there is no undo log and no purge or VACUUM step.

MySQL / InnoDBPostgreSQLMongoDB / WiredTiger
Where rows / documents liveClustered B+ tree on the primary keyHeap; every index points to a TIDB-tree keyed by RecordId; every index, _id too, points to it
UPDATEIn place; old version to the undo logNew tuple; xmax on the old oneNew version on the in-memory update chain
DELETEDelete-mark; purge removes itSet xmax; VACUUM removes itTombstone; reconciliation leaves it out
Old versions for readersUndo logDead tuples in the heapUpdate chains, then the history store
Crash safetyRedo log (+ doublewrite)WAL (+ full-page images)Journal + copy-on-write checkpoints
Replication logBinlog (2PC with redo)The WAL itselfThe oplog, a collection written in the same transaction

insertOne, updateOne and deleteOne, step by step

Every write runs in one WiredTiger transaction. The transaction changes the collection, every index and the oplog, and all of it becomes visible at the commit timestamp.

  • insertOne: the document gets the next RecordId and goes into the collection table, then one key goes into each index. A key already in the unique _id_ index gives E11000 duplicate key error. The transaction is rolled back and its uncommitted update simply dropped (Demo: duplicate _id).
  • updateOne: the document is found through _id_, the update is applied, and a new version goes onto the update chain. For every index whose key changes, the old key gets a tombstone and the new key is inserted. If $set changes nothing, nothing is written: modifiedCount: 0 and no oplog entry (Demo: updateOne).
  • deleteOne: a tombstone goes on the document and on each of its index keys (Demo: deleteOne, then a checkpoint).

The journal, checkpoints and crash recovery

At commit, the transaction's log record goes into the journal buffer. The journal (WiredTiger's write-ahead log) is written and fsync()ed in two cases: every storage.journal.commitIntervalMs (100 ms), and at once when a write asks for j: true (which w: "majority" implies).

The .wt data files are only brought up to date by a checkpoint, every 60 s by default (storage.syncPeriodSecs). A checkpoint writes new page images to new places in the files (copy-on-write), so the previous checkpoint stays valid until the new one is complete. In a replica set, a checkpoint contains only data up to the stable timestamp, which is the majority commit point.

After a crash, WiredTiger opens the last checkpoint and replays the journal written after it (Demo: checkpoint, crash, journal replay).

The oplog and replication

Each write also inserts an entry into local.oplog.rs, in the same transaction. The oplog is a capped collection, so the oldest entries are dropped as new ones arrive. Entries are idempotent, so applying one twice gives the same result:

{ts: t11, op: "i", ns: "test.users", o: {_id: 4, name: "Dan", age: 31}}
{ts: t12, op: "u", ns: "test.users", o: {$v: 2, diff: {u: {age: 24}}}, o2: {_id: 5}}
{ts: t13, op: "d", ns: "test.users", o: {_id: 3}}

Each secondary keeps a getMore open on its sync source's oplog (a tailable cursor). It applies the batches it gets in parallel and journals them. Then it reports its position (replSetUpdatePosition). The primary moves the majority commit point once a majority of members have an entry durable.

Write concern, elections and rollback

writeConcernThe ack waits forIf the primary dies right after the ack
{w: "majority"} (default since 5.0)The journal of a majority of members (j: true is implied)Safe. Any new primary must have the write, because a candidate only gets votes if its oplog is at least as recent as the voter's.
{w: 1}The primary's in-memory commitCan be lost. If no secondary copied it, the new primary does not have it, and the old primary rolls it back when it rejoins.
{w: 1, j: true}The primary's journalSurvives a restart of that server, but can still be rolled back in a failover.
Timeline with lanes client, A, B and C: insertOne reaches A, which commits in memory and sends the w: 1 ack at once; B and C apply and journal the write, B reports its position, and only then does A send the w: majority ack.
w: 1 acknowledges as soon as the primary has the write; w: majority waits until a majority has journaled it, so a failover cannot roll it back.

Members send heartbeats every 2 s. If the primary is silent for electionTimeoutMillis (10 s), a secondary starts an election for a new term and asks for votes (replSetRequestVotes). The old primary comes back and finds that its oplog has entries the new primary lacks. It enters ROLLBACK: it rolls its data back to the stable timestamp, saves the undone documents in rollback/ files, and catches up. A member with a higher priority then calls a priority-takeover election (Demo: w: 1 write lost in a failover). Reads have a similar choice: readConcern: "majority" returns only data at or before the majority commit point. The default, "local", returns the newest data, which could still be rolled back.

See also How MySQL Runs a Query, How PostgreSQL Runs a Query, MVCC and Isolation Levels, LSM Tree vs B+ Tree, Database Replication and Raft (MongoDB's election protocol is based on it).

What the page leaves out

  • Sharding: mongos, config servers, chunks and balancing.
  • Multi-document transactions and write conflicts between clients.
  • Eviction under cache pressure and the history store.
  • Real page splits. When a leaf fills up, the page runs a checkpoint and rebalances the tree instead.
  • Compression and prefix compression of index keys.
  • The oplog's own table and its truncation.
  • Change streams, read preference (reading from secondaries), readConcern: "snapshot", and retryable-write bookkeeping.
  • After a crash, A really stays a secondary until an election. Here it is re-elected at once.