Change Streams: Reacting to Data in Real Time

Advanced
12 min

Change Streams: Reacting to Data in Real Time

Applications often need to know when data changes: invalidate a cache, push a notification, sync a search index, or feed an event pipeline. Polling with timestamps is fragile and slow. Change streams let a client subscribe to a live feed of inserts, updates, replaces and deletes, backed by the replica set oplog. After this lesson you will be able to open a change stream at collection, database or deployment level, filter and shape its events, see full documents on updates, and resume reliably after a restart.

Opening a Stream

watch() exists on collections, databases and the client. It returns a cursor that stays open and yields one event per committed change:

javascript
const orders = client.db("shop").collection("orders"); const stream = orders.watch(); // one collection const dbStream = client.db("shop").watch(); // every collection in shop const allStream = client.watch(); // whole deployment except admin/local/config stream.on("change", (event) => console.log(event)); // event-emitter style // or: for await (const event of stream) { ... }

Change streams require a replica set or sharded cluster and only report majority-committed changes, so an event can never be rolled back after you have seen it.

Anatomy of a Change Event

javascript
{ _id: { _data: "8265F1..." }, // resume token operationType: "update", // insert | update | replace | delete | drop | rename | invalidate ... clusterTime: Timestamp({ t: 1710000000, i: 3 }), wallTime: ISODate("2026-03-09T10:00:00Z"), ns: { db: "shop", coll: "orders" }, documentKey: { _id: ObjectId("...") }, updateDescription: { updatedFields: { status: "shipped" }, removedFields: [], truncatedArrays: [] }, fullDocument: { _id: ObjectId("..."), status: "shipped", ... } // only with fullDocument option }

For inserts and replaces fullDocument is always present. For updates you receive only updateDescription unless you ask for more:

| Option | Effect | |---|---| | fullDocument: "updateLookup" | fetch the current document after the update (may include later changes) | | fullDocument: "whenAvailable" / "required" | the exact post-image, when pre/post images are enabled | | fullDocumentBeforeChange: "whenAvailable" / "required" | the pre-image; the only way to see what a delete removed |

Pre- and post-images are stored per collection and must be switched on:

javascript
db.runCommand({ collMod: "orders", changeStreamPreAndPostImages: { enabled: true } })

Filtering and Shaping Events

The first argument to watch() is an aggregation pipeline that runs on the server, so unwanted events never cross the network. Allowed stages are $match, $project, $addFields/$set, $unset, $replaceRoot/$replaceWith and $redact.

javascript
orders.watch([ { $match: { $or: [ { operationType: "insert" }, { operationType: "update", "updateDescription.updatedFields.status": "shipped" } ] } }, { $project: { operationType: 1, documentKey: 1, "fullDocument.customerId": 1, "fullDocument.status": 1 } } ], { fullDocument: "updateLookup" })

Resuming After a Restart

Every event's _id is a resume token. Store it after you finish processing each event; when your process restarts, pass it back:

javascript
const token = await loadResumeToken(); const stream = orders.watch([], token ? { resumeAfter: token } : {});
  • resumeAfter continues after the given event.
  • startAfter does the same but also works after an invalidate event (collection dropped or renamed).
  • startAtOperationTime starts from a cluster timestamp instead of a token.

Resuming only works while the token's position is still in the oplog, so size the oplog for the longest downtime you expect. Because you persist the token after processing, delivery is at-least-once: make handlers idempotent.

Where Change Streams Fit

Typical uses: cache invalidation, real-time dashboards over WebSockets, syncing documents into Elasticsearch or Atlas Search, audit trails, and event-driven microservices. Keep the handler fast — a slow consumer builds a backlog in the cursor and eventually falls behind the oplog window.

Common Mistakes

  • Not persisting the resume token, so a restart silently skips every change made while the process was down.
  • Expecting fullDocument on updates without fullDocument: "updateLookup" or post-images.
  • Running against a standalone server, which rejects watch() with "The $changeStream stage is only supported on replica sets".
  • Doing heavy synchronous work inside the loop, which stalls event delivery.
Quick Quiz
Question 1 of 3

Which option makes update events include the whole current document?

Key Takeaways

  • watch() on a collection, database or client returns a live cursor of majority-committed change events; a replica set is required.
  • Events carry operationType, documentKey, updateDescription and, when requested, fullDocument and fullDocumentBeforeChange.
  • A pipeline passed to watch() filters and projects events on the server.
  • Persist each event's _id and resume with resumeAfter or startAfter; handlers must be idempotent.
  • Use change streams for cache invalidation, notifications, search-index sync and event-driven integrations.

Next lesson: Geospatial Queries with 2dsphere Indexes — store locations as GeoJSON and query by distance, area and intersection.

Change Streams: Reacting to Data in Real Time - MongoDB | CodeYourCraft | CodeYourCraft