MongoDB Change Streams
Applications often need to react the moment data changes. A shop wants to send an alert when a new order arrives, and a dashboard wants to refresh when a stock level drops. Change streams let an application listen to a collection and receive each change as it happens, without checking the database again and again.
The Doorbell Picture
Checking a collection every few seconds works like walking to the front door every minute to see whether a visitor has arrived. A change stream works like a doorbell. The database rings the application only when something happens.
Polling (old way) Change Stream (new way)
App ---> "Any change?" ---> DB App <--- "New order!" <--- DB
App ---> "Any change?" ---> DB App <--- "Price updated!" <--- DB
App ---> "Any change?" ---> DB
(many wasted requests) (messages arrive only on change)Requirements
- Change streams need a replica set or a sharded cluster. A single standalone server cannot run them.
- A local practice setup can use a one-member replica set.
- The user account needs permission to read the collection.
Open a Change Stream
The watch() method opens a stream on a collection, a database, or a whole deployment. The shell example below listens to an orders collection:
const stream = db.orders.watch();
while (!stream.isClosed()) {
const change = stream.tryNext();
if (change !== null) {
printjson(change);
}
}Open a second shell window and insert a document to see the stream react:
db.orders.insertOne({ item: "Pen", qty: 10 })What a Change Event Looks Like
{
operationType: "insert",
ns: { db: "shop", coll: "orders" },
documentKey: { _id: ObjectId("...") },
fullDocument: { _id: ObjectId("..."), item: "Pen", qty: 10 }
}| Field | Meaning |
|---|---|
| operationType | The kind of change, such as insert, update, replace, or delete |
| ns | The database and collection where the change happened |
| documentKey | The _id of the changed document |
| fullDocument | The complete new document, present for inserts |
| updateDescription | The fields that changed, present for updates |
Operation Types
- insert fires when a document is added.
- update fires when fields change inside a document.
- replace fires when a whole document is swapped.
- delete fires when a document is removed.
- drop and rename fire when a collection is dropped or renamed.
Get the Full Document for Updates
An update event normally lists only the changed fields. Ask MongoDB to attach the complete current document with the fullDocument option:
const stream = db.orders.watch([], { fullDocument: "updateLookup" });Each update event now carries the whole document as it looks after the change.
Filter the Stream
A pipeline narrows the events that reach the application. The filter below keeps only inserts for large orders:
const pipeline = [
{ $match: { operationType: "insert", "fullDocument.qty": { $gte: 100 } } }
];
const stream = db.orders.watch(pipeline);Filtering on the server saves network traffic and processing time in the application.
Resume After a Break
Every event includes a resume token in its _id field. Think of it as a bookmark in a book. An application that stops, crashes, or restarts saves the last token and resumes from that exact spot:
const stream = db.orders.watch([], { resumeAfter: savedToken });The stream replays every change that happened while the application was away. The replay works only while the changes remain in the replication log, so a very long outage can exceed that window.
Resume Flow
Event 1 (token A) ---> processed, token A saved
Event 2 (token B) ---> processed, token B saved
... application crashes ...
Event 3 (token C) ---> happens while app is offline
... application restarts ...
resumeAfter: token B ---> receives Event 3 and continuesReal-World Uses
- Send an email or message when a new order arrives.
- Refresh a live dashboard with new sales numbers.
- Copy changes into a search engine or another database.
- Write an audit trail that records who changed what.
- Clear a cache when the source data changes.
Usage in Node.js
const changeStream = client.db("shop").collection("orders").watch();
changeStream.on("change", function (event) {
console.log("Change detected:", event.operationType);
});Good Practices
- Store the latest resume token in a safe place after each processed event.
- Keep the processing work inside the listener short and fast.
- Close the stream when the application shuts down.
- Filter events on the server to receive only the changes that matter.
Summary
Change streams push database changes to applications in real time. The watch() method opens a stream on a replica set or sharded cluster. Each event describes the operation type, the document key, and optionally the full document. Pipelines filter events on the server, and resume tokens let an application continue after a break without missing any change.
