Data Warehouse Insider

MongoDB Change Stream Cursor Internals

Change streams rely on MongoDB's oplog and tailable cursors, not a separate feed.

Staff Writer · · 10 min read
Cover illustration for “MongoDB Change Stream Cursor Internals”
Database Replication Internals · October 6, 2026 · 10 min read · 2,308 words

The oplog comes first. Everything a change stream does, and everything it cannot do, traces back to a single structure inside MongoDB's replication system: the operations log, or oplog, a capped collection that records every write a replica-set primary accepts, in durable, ordered sequence. It logs operations, not states. An insert, an update, a delete each get written as an entry; the result of a read query never does, because there is nothing for replication to propagate when nothing changes. Secondary nodes read this same log to stay in sync with the primary, applying each entry in order to reproduce the primary's state locally. Change streams do not have their own separate feed of database activity. They read the oplog, the identical structure replication itself depends on.

This shared dependency explains a constraint that trips up anyone moving from a standalone MongoDB instance to a replica-set deployment for the first time: change streams require a replica set or sharded cluster, because the oplog only exists when replication is active. A standalone server has no replication to perform, so it keeps no oplog, and with no oplog there is nothing for a change stream to read. For local development, a single-node replica set, started with a command like mongod --replSet rs0, satisfies this requirement without needing multiple servers. Atlas sidesteps the question entirely, since every Atlas cluster is provisioned as a replica set or a sharded cluster, where each shard is itself a replica set with its own oplog. The oplog's capped, circular nature (what gets overwritten as it fills, and why that overwriting is the source of nearly every hard limit practitioners run into) is the subject of a later section. What matters here is the ordering of ideas: the oplog is the primitive, and the change stream is built on top of it.

The tailable cursor as the mechanical primitive beneath every change stream

A capped collection is only useful for real-time consumption if something can read it as it grows, without re-querying from the start each time. That something is the tailable cursor, the primitive that makes it possible to follow a capped collection as new entries arrive, and it is the mechanism every change stream sits on top of. An ordinary MongoDB cursor closes once it exhausts its result set. A tailable cursor behaves differently: when it reaches the end of the available data, it stays open rather than closing, and resumes delivering results the moment new documents are appended to the collection. Before change streams existed as an API, real-time reactivity in MongoDB meant tailing the oplog directly with a cursor of this kind, a technique that worked but demanded a fair amount of custom engineering to implement and keep running correctly.

The practical mechanics matter. Setting tailable: true together with awaitData: true and a maxAwaitTimeMS value produces long-polling behavior: the server holds the connection open, up to the specified timeout, waiting for new oplog entries to appear, and only returns once something new arrives or the timeout expires. This is what keeps a tailing consumer from hammering the server with a tight, repeated polling loop. It is a deliberately patient cursor, built to wait.

Direct oplog tailing predates change streams, and nothing prevents it today. A consumer can query local.oplog.rs directly, opening a cursor with { tailable: true, awaitData: true } and filtering by the ts field and the operation type recorded in each entry (i for insert, u for update, d for delete, c for command). The ts field is a BSON Timestamp, a compact structure encoding two 32-bit integers: seconds since a reference epoch, and an ordinal value that disambiguates multiple entries occurring within the same second. A consumer tailing the oplog directly takes on real responsibility for durability. To resume cleanly after a restart, it has to save the last processed ts value to durable storage, then reopen its cursor with a filter like ts: { $gt: lastTs } to pick up exactly where it left off. A primary failover complicates this further: the consumer must detect that a new primary has been elected, reconnect to it, and resume from the same ts value, and all of that detection and reconnection logic belongs entirely to the consumer. MongoDB provides none of it automatically.

There is a deeper problem with building on the oplog entry format directly. It is an internal structure of MongoDB's replication system, not a documented, stable public API, and its layout has changed across MongoDB versions. Code written against the oplog's shape in one release is not guaranteed to work against the next. Direct tailing carries real advantages: lower-level access to raw operation types and update diffs, lower overhead for simple streaming use cases, and compatibility with older MongoDB deployments that predate the change stream API. It also carries real costs: no built-in resume tokens, no pre-image support, no unified view across a sharded cluster, and no access at all on Atlas, where the local database is blocked from client queries. The tailable cursor is a real, working mechanism, not a historical curiosity. Change streams did not replace it; they wrapped it.

What the change stream API adds to the tailable cursor

Change streams are not a different mechanism from oplog tailing. They are a defined, versioned API surface built on top of the same tailable cursor and the same oplog, adding three things a raw tailing implementation has to build by hand: event normalization, aggregation pipeline filtering, and a resume token system. Those three additions are what make oplog consumption usable by application developers who have no need to understand MongoDB's internal replication format.

The watch() method is the entry point, and it can be called on a single collection, an entire database, or a full deployment. Change streams, introduced in MongoDB 3.6, started at the collection level. MongoDB 4.0 extended the capability to an entire database or full deployment, made the fullDocument: "updateLookup" option available at that broader scope (it had already been available at the collection level since 3.6), and added sharded cluster support without the restrictions that earlier versions carried. MongoDB 6.0 added pre-image support through fullDocumentBeforeChange.

Every event a change stream delivers arrives in a normalized, version-stable schema: operationType, documentKey, clusterTime, fullDocument where applicable, updateDescription for update operations, and ns identifying the namespace. One detail catches developers off guard regularly: update events, by default, carry only what changed, the updatedFields and removedFields inside updateDescription, not the full document as it now stands. Anyone expecting the complete document on every update needs to request it explicitly. Passing fullDocument: "updateLookup" tells MongoDB to fetch the current document at the time the event is read. That distinction has a real consequence: if the document was modified again between the original change and the lookup, the consumer receives the later state, not the state immediately following the original write, and if the document was deleted in the meantime, the consumer receives null.

Change streams accept an aggregation pipeline passed into watch(), which lets filtering and projection happen on the server, before events ever reach the consumer. A $match stage can restrict a stream to insert operations only, for instance; other supported stages include $project, $addFields, $replaceRoot, and $redact. This is not a cosmetic convenience. Every open change stream cursor makes the server scan the oplog regardless of how much of that activity the consumer actually wants, so shipping events the consumer immediately discards wastes both server-side work and network bandwidth. Filtering early, inside the pipeline, is a meaningful performance decision in any deployment handling real write volume, not an optional refinement. One constraint is strict and non-negotiable: projecting away the _id field of the event document causes an error. That field carries the resume token, and MongoDB will not produce an event it cannot offer a resumption point for.

Pre-image support, available from MongoDB 6.0 through fullDocumentBeforeChange, stores the state of a document as it existed before an update or delete, writing it into config.system.preimages. It is the only mechanism that reveals what a deleted document actually contained, since a delete event on its own carries no content. Pre-images are not enabled by default. Each collection needs changeStreamPreAndPostImages: { enabled: true } set explicitly, and because storage consumption scales with write rate, the sensible approach is to set an expireAfterSeconds value and turn the feature on only where the application genuinely needs that history.

How resume tokens work and what guarantees they provide

Resume tokens make change streams restartable after a disconnect, a crash, or a planned deployment, but a resume token is only a pointer into the oplog, as durable as the oplog entry it points to. Every change stream event carries a resume token in its _id field, marking a specific position in the oplog. Passing that token back in as resumeAfter when opening a new stream restarts delivery from precisely that position. When no prior token exists, startAtOperationTime offers an alternative entry point: instead of a token tied to one specific event, it accepts a cluster timestamp and begins delivery from there.

The delivery guarantee a consumer actually gets is at-least-once, not exactly-once, and the distinction has direct consequences for how a consumer has to be written. Consider the two ways a consumer can order its work relative to saving the resume token. If it saves the token first and then confirms delivery to its downstream destination, a crash between those two steps loses the event permanently: the token has already advanced past it. If it confirms delivery first and then saves the token, a crash during the save means the event gets reprocessed the next time the consumer starts up, since the stored token still points to the position before that event. The second failure mode is the survivable one. The correct pattern persists the resume token only after the corresponding event has been durably delivered to its destination. Consumer handlers have to be written as idempotent operations, built to tolerate receiving the same event more than once without corrupting downstream state.

The resume token's usefulness is bounded by the same structure that bounds everything else about the oplog: its fixed size. Once the oplog has cycled past the entry a given token points to, that token can no longer be used to resume, full stop, and the consumer has no way to recover the gap in between. Under heavy write load, the window of time the oplog actually covers can shrink to a matter of hours, or less. Checking the current window with rs.printReplicationInfo() and sizing the oplog around worst-case downtime, not average downtime, is the only defense against this failure. When a consumer's downtime exceeds that window anyway, the only path forward is a full re-snapshot of the data, an expensive and disruptive operation at any meaningful scale. This is the point where a change stream stops behaving like a convenient API and starts behaving like a system with real, specific failure modes that have to be designed around deliberately.

Diagram: At-Least-Once Delivery: The Two Failure Modes. Visualizes: Show two contrasting sequences illustrating the at-least-once delivery guarantee for change stream consumers.

The operational constraints that follow directly from the oplog's capped-collection design

Every hard limit a production change stream deployment runs into traces back to the same fact: the oplog is a capped collection, served through server-side cursors, and nothing about that architecture changes once an application starts depending on it at scale.

The oplog window is the single biggest operational constraint for any change-data-capture workload built on change streams. Because the oplog is a fixed-size circular buffer, high write volume can compress its coverage down to a matter of hours. A consumer down longer than that window cannot resume from its saved token; it has to re-snapshot from scratch. The oplog should be sized for the worst realistic consumer downtime an operations team can anticipate, not the average case, and rs.printReplicationInfo() gives a direct read on the current window at any moment.

Cursor count is a real server resource, not a free abstraction, because each open change stream is a server-side cursor actively scanning the oplog. Tens of concurrent streams run without issue; thousands do not. When multiple consumers need the same set of events, the architecturally sound pattern fans those events out through a single upstream consumer rather than opening a dedicated change stream per client. Part of why cursor count matters so much is that the oplog collection supports no indexes, so every cursor reading it does so sequentially, and the aggregate cost of that scanning grows directly with how many cursors are open at once.

Sharded clusters introduce a latency source that single replica sets never encounter. A mongos router has to poll every shard, including shards with little or no recent write activity, in order to maintain a globally ordered view of change events across the cluster. A shard that has not recently written a noop heartbeat entry stalls that ordering for the whole deployment. Lowering periodicNoopIntervalSecs on cold shards reduces this latency, at the cost of a higher volume of oplog writes on those shards. This tradeoff exists because change streams offer something genuinely valuable: a globally ordered view of changes across every shard in a cluster, built on a global logical clock, and that ordering guarantee is what makes polling the cold shards necessary.

Change stream event documents are bound by the 16 MB BSON document limit, the same ceiling that applies to any document in MongoDB's general data model. A notification document that exceeds that size, which can happen when a full pre-image or post-image is attached to a large document, fails to deliver. From MongoDB 6.0.9 onward, the $changeStreamSplitLargeEvent aggregation stage addresses this directly, splitting an oversized event into multiple deliverable pieces. Each of these constraints, the oplog window, the cursor ceiling, the cold-shard latency, the 16 MB limit, follows from decisions made deep in the replication layer, long before the change stream API existed to expose them to application developers.

Sources

  1. An Introduction to Change Streams
  2. Change Streams Production Recommendations - Database Manual - MongoDB Docs
  3. MongoDB Change Streams - Database Manual - MongoDB Docs
  4. Motor Tailable Cursor Example - Motor 3.7.1 documentation
  5. Using MongoDB 3.6 Change Streams - Percona
  6. Change stream permanent oplog - Working with Data - MongoDB Community Hub
  7. How Do Change Streams Work In MongoDB?
  8. Efficiency of change streams when large oplog entries are created - Working with Data - MongoDB Community Hub

More in Database Replication Internals