Most MongoDB problems in production are not about the document model. They are about what happens underneath it: a working set that outgrew the cache, a query that picked the wrong index, a write acknowledged by one node and rolled back after a failover, or a query without the shard key that quietly visits every shard. Whether MongoDB fits your data at all is a separate question, covered in MongoDB: when it actually fits. This article assumes you run it and want to understand the engine well enough to predict its behaviour.

We follow one operation down the stack: from the driver, through mongos in a sharded cluster, into a shard's primary, where the query layer chooses a plan, the replication layer records an oplog entry and the WiredTiger storage engine holds the data in cache and makes it durable. At each layer you will see the mechanism, the metric that exposes it and the failure it causes when ignored.

Advertisement

The stack in one picture

Every mongod process has the same three layers. The query layer parses queries and aggregation pipelines, chooses an index and executes the plan. The replication layer records each change in the oplog and manages elections. The storage engine, WiredTiger, stores collections and indexes as B-trees, caches pages in memory and writes them to disk. A sharded cluster adds mongos routers, which hold a cached copy of which shard owns which range of shard-key values, and a config server replica set that stores that metadata.

One sharded MongoDB read and write, layer by layerDriverpool, retries, concernsmongoscached routing tableConfig serversreplica set: rangesrefreshShard A replica set: primaryQuery layerplanner, plan cache, aggregationReplicationoplog entry in the same storage transactionWiredTigerB-tree pages in cache, MVCC snapshotsjournal (write-ahead) + checkpoint every 60 stargetedSecondary 1fetch oplog, applySecondary 2fetch oplog, applyShard B, C ...only if the querylacks the shard keypullpullscatter-gatherMajority commit pointw: majority acks, read concern majorityWrites become durable on one node through the journal and safe against failover through the majority commit point.Reads are cheap only when the shard key routes them and an index matches their filter and sort.
A targeted read goes from the driver through mongos to one shard; a query without the shard key fans out to every shard. Inside a shard, the primary records each change in its oplog and secondaries pull and apply it.

WiredTiger: cache, snapshots and durability

WiredTiger keeps an internal cache of uncompressed B-tree pages. By default it takes the larger of 50% of (RAM minus 1 GB) or 256 MB, so a 64 GB server gets 31.5 GB. On disk, collection blocks are compressed with snappy by default and indexes use prefix compression, so the rest of the memory is not wasted: the operating system's file cache holds compressed blocks, which are cheaper to bring into the WiredTiger cache than a disk read. The working set, the documents and index pages your queries actually touch, should fit in the WiredTiger cache. When it does not, reads go to disk and latency follows the disk.

Concurrency uses multi-version concurrency control. Each operation reads from a point-in-time snapshot, and writers to different documents in the same collection do not block each other; only intent locks are taken at the global, database and collection levels. When two operations modify the same document, one gets a write conflict and MongoDB retries it transparently, except inside a multi-document transaction, where the error is returned for the client to retry. The cost of MVCC is old versions: a long-running transaction or snapshot read pins history, which takes cache space. That is one reason transactions are capped at 60 seconds by default.

Durability combines two mechanisms. A checkpoint writes a consistent snapshot of all data to disk every 60 seconds. Between checkpoints, the journal is a write-ahead log of every modification; after a crash, MongoDB loads the last checkpoint and replays the journal. A write with j: true is acknowledged only after its journal record is on disk.

When the cache fills with pages, or with modified pages not yet written, background eviction threads push them out. If they cannot keep up, application threads are drafted into eviction, and every operation slows at once. The signature is a latency spike across all collections together with high dirty-cache figures in db.serverStatus().wiredTiger.cache. The usual causes are a working set that outgrew memory, a bulk load or a long-running snapshot pinning history.

Advertisement

The query planner and reading explain

For a new query shape (the combination of filter fields, sort and projection), the planner lists the candidate indexes, builds a plan for each and runs them against each other for a short trial. The plan that produces results with the least work wins and is stored in the plan cache for that shape. Later queries of the same shape reuse it. If a cached plan starts doing far more work than it did when cached, the planner re-evaluates it; the cache is also cleared when indexes change or the process restarts.

This is efficient, but it explains a classic incident: a plan chosen during a trial on typical values is reused for a skewed value, such as a customer with ten million orders, and does far more work there. Diagnose with explain("executionStats") and compare three numbers: documents returned, index keys examined and documents examined. A healthy plan examines close to as many keys as it returns. A COLLSCAN stage means no index was used; a SORT stage above an IXSCAN means the sort happens in memory, which is limited to 100 MB per stage unless it can spill to disk.

The rule that prevents most bad plans is Equality, Sort, Range: in a compound index, put fields tested for equality first, then the sort fields, then fields tested with ranges. With that order the index already returns documents in sort order for each equality prefix, so limit(20) can stop after about 20 keys. General B-tree behaviour is covered in B-tree indexes.

// Filter on customerId (equality) and total (range), sort by createdAt.
const q = { customerId: 42, total: { $gt: 100 } };

// Index A puts the range before the sort: the planner must sort in memory.
db.orders.createIndex({ customerId: 1, total: 1, createdAt: -1 });
// Index B follows Equality, Sort, Range: index order already matches the sort.
db.orders.createIndex({ customerId: 1, createdAt: -1, total: 1 });

const e = db.orders.find(q).sort({ createdAt: -1 }).limit(20).explain("executionStats");
printjson({
  plan: e.queryPlanner.winningPlan,              // look for SORT above IXSCAN, or COLLSCAN
  nReturned: e.executionStats.nReturned,
  keys: e.executionStats.totalKeysExamined,
  docs: e.executionStats.totalDocsExamined,
  ms: e.executionStats.executionTimeMillis,
});

Worked example: an orders page

An orders collection holds 40 million documents. The account page runs the query above for a customer with 12,000 orders, 3,000 of them over 100. With index A, the plan examines 3,000 keys, fetches 3,000 documents, sorts them in memory and returns 20: about 150 times more work than needed, growing with the customer's history. With index B, the plan walks the customer's orders newest first, checks total from the index key without fetching, and stops at the twentieth match: around 80 keys and 20 documents. The latency difference is small for a typical customer and large for the biggest ones, which is exactly the distribution that pages someone. After adding index B, drop index A unless another query needs it: every index adds work to every write and occupies cache.

Replication: the oplog, elections and rollback

The primary records each change in the oplog, a capped collection in the local database, in the same storage transaction as the change itself. Operations are rewritten into idempotent form, so an increment becomes a set to the resulting value and an update to many documents becomes one entry per document. Secondaries pull batches from a sync source, which may be another secondary, and apply them in parallel. Generic replication theory is in database replication.

The oplog's size defines the replication window: how long a secondary can be offline and still catch up. By default it is 5% of free disk space, at least 990 MB and at most 50 GB, and you can set a minimum retention in hours. A heavy bulk update rewrites one entry per document and can shrink the window from days to minutes. A secondary that falls off the end needs a full initial sync. Monitor the window and replication lag, not just lag.

Members exchange heartbeats every two seconds. If a secondary cannot reach the primary for the election timeout (10 seconds by default), it calls an election with a higher term, and a candidate needs votes from a majority of voting members. MongoDB's documentation says the median time to elect a new primary should typically not exceed 12 seconds with default settings; drivers retry eligible writes once across the change.

The majority commit point is the newest oplog entry replicated to a majority of voting members. A write acknowledged with w: "majority" is behind that point and survives any failover. A write acknowledged with w: 1 may exist only on the old primary; if that node rejoins after a new primary was elected without it, the write is rolled back and saved to rollback files on disk, not replayed. Read concern majority reads only data behind the commit point, so it never shows a value that can be rolled back.

from pymongo import MongoClient, WriteConcern
from pymongo.read_concern import ReadConcern

client = MongoClient("mongodb://db1,db2,db3/?replicaSet=rs0&retryWrites=true")
orders, stock = client.shop.orders, client.shop.stock

def place_order(session, order):
    # Both writes commit or neither does. A WriteConflict inside the callback surfaces
    # as TransientTransactionError, and with_transaction re-runs the whole callback.
    res = stock.update_one(
        {"_id": order["sku"], "qty": {"$gte": order["qty"]}},
        {"$inc": {"qty": -order["qty"]}},
        session=session,
    )
    if res.modified_count != 1:
        raise ValueError("insufficient stock")     # aborts; not retried
    orders.insert_one(order, session=session)

with client.start_session() as s:
    s.with_transaction(
        lambda sess: place_order(sess, {"sku": "A-1", "qty": 2, "customerId": 42}),
        read_concern=ReadConcern("snapshot"),
        write_concern=WriteConcern("majority"),
    )

Sharding: routing, the balancer and orphans

mongos routes using its cached routing table. A query whose filter includes the shard key, or a prefix of a compound one, goes only to the shards owning the matching ranges. A query without it is scatter-gather: sent to every shard, with results merged on mongos. Scatter-gather is correct but its latency is the slowest shard's, and its load grows with the number of shards, so adding shards makes it worse. Choosing the key itself is covered in the fit article and in database sharding.

The balancer moves ranges of data between shards. It acts per collection when the data size difference between the fullest and emptiest shard reaches three times the range size; with the default range size of 128 MB, that is 384 MB. A migration copies the documents, catches up on writes made during the copy, takes a brief critical section to commit the new ownership in the config servers, and then the donor deletes its copy in the background. Until deletion, those documents are orphans: primary reads filter them out, but reads with read concern available on secondaries can return them. A range whose documents all share one shard-key value cannot be split and becomes a jumbo range that the balancer cannot move.

Failure modes

  • Working set outgrows the cache. Disk reads rise, then eviction pressure drafts application threads. Watch cache usage, dirty bytes and pages read into cache per second.
  • Wrong cached plan. One query shape suddenly slows for some values. Compare keys and documents examined; fix the index order rather than relying on hint().
  • Rollback after failover. Writes acknowledged with w: 1 disappear from the new primary. Use majority writes for anything you cannot afford to lose.
  • Oplog window collapse. A bulk update shrinks the window and a lagging secondary needs a full resync. Batch and throttle bulk updates, and size the oplog for maintenance windows.
  • Scatter-gather growth. A frequent query without the shard key gets slower with every shard added.
  • Long transactions. They pin history in cache, raise conflict rates and hit the 60-second limit. Keep transactions short and touch few documents.

Operational trade-offs

DecisionCheaper optionSafer or faster optionCost of the safer option
Write concernw: 1w: majorityOne extra replication round trip per write
Read concernlocalmajorityMay return slightly older data
IndexesFew, broadOne per hot query shapeWrite amplification and cache use
Oplog sizeDefaultSized for the longest maintenanceDisk space
Shard keyMonotonic, simpleMatches the hottest query filterHarder to change later

What to do next

  1. Measure the working set against the WiredTiger cache and track dirty cache and eviction metrics on every primary.
  2. Run explain("executionStats") on your ten most frequent query shapes and flag any with COLLSCAN, in-memory SORT or keys examined far above documents returned.
  3. Reorder compound indexes to Equality, Sort, Range, and drop indexes no query uses.
  4. Set w: "majority" for writes you cannot lose and check that nothing overrides it per operation.
  5. Record the oplog window in hours and alert when it drops below your longest planned maintenance.
  6. In sharded clusters, list queries that lack the shard key and estimate their cost at twice the current shard count.
  7. Rehearse a primary failover in staging and confirm the application retries and loses nothing. For storage engine background, compare with LSM trees.
Key takeaway: MongoDB is a query layer, a replication layer and the WiredTiger storage engine. WiredTiger caches uncompressed B-tree pages, isolates operations with MVCC snapshots and makes writes durable through the journal and 60-second checkpoints. The planner races candidate plans and caches the winner per query shape, so index order following Equality, Sort, Range matters. Replication flows through an idempotent oplog whose size is the recovery window, and only majority-acknowledged writes survive failover. In sharded clusters, queries without the shard key visit every shard, and the balancer moves data by size difference between shards.