Skip to content

Change Streams

Change streams let you watch for data changes in real time — like a live feed of every insert, update, and delete.

Think of a newspaper subscription:

  • Without change streams: You call the newspaper every hour to ask “Any news?” (polling)
  • With change streams: The newspaper delivers new editions to your doorstep the moment they’re printed (push notifications)
sequenceDiagram
participant App as Application
participant MongoDB as MongoDB
App->>MongoDB: db.collection.watch()
MongoDB-->>App: ✅ Change stream opened
Note over App,MongoDB: Meanwhile, other apps modify data...
OtherApp->>MongoDB: INSERT { name: "Alice" }
MongoDB-->>App: { operationType: "insert", fullDocument: {...} }
OtherApp->>MongoDB: UPDATE { $set: { age: 30 } }
MongoDB-->>App: { operationType: "update", updateDescription: {...} }
OtherApp->>MongoDB: DELETE { _id: ObjectId(...) }
MongoDB-->>App: { operationType: "delete", documentKey: {...} }
// Watch all changes on a collection
const changeStream = db.users.watch();
// Listen for changes
changeStream.on('change', (change) => {
console.log('Change detected:', change);
});
// Output:
// { operationType: 'insert', fullDocument: { name: 'Alice', ... } }
// { operationType: 'update', updateDescription: { updatedFields: { age: 30 } } }
// { operationType: 'delete', documentKey: { _id: ObjectId(...) } }
const User = require('./models/User');
async function watchUsers() {
const changeStream = User.watch();
changeStream.on('change', (change) => {
switch (change.operationType) {
case 'insert':
console.log('🆕 New user:', change.fullDocument.email);
// Send welcome email, update cache, notify admins
break;
case 'update':
console.log('✏️ User updated:', change.documentKey._id);
// Invalidate cache, log the change
break;
case 'delete':
console.log('🗑️ User deleted:', change.documentKey._id);
// Clean up related data, notify admins
break;
}
});
console.log('👀 Watching for user changes...');
}
// Watch only inserts and updates
const changeStream = db.users.watch([
{ $match: { operationType: { $in: ['insert', 'update'] } } }
]);
// Watch changes to specific fields
const changeStream = db.orders.watch([
{ $match: { 'updateDescription.updatedFields.status': { $exists: true } } }
]);
// Watch changes on a specific document
const changeStream = db.users.watch([
{ $match: { 'documentKey._id': ObjectId('user_id_here') } }
]);

Change streams support resume tokens — if your app crashes, you can resume from where you left off.

// Save the resume token
let resumeToken = null;
const changeStream = db.users.watch();
changeStream.on('change', (change) => {
resumeToken = change._id; // save this token
processChange(change);
});
// On restart, resume from the saved token
const changeStream = db.users.watch([], { resumeAfter: resumeToken });
Use CaseHow Change Streams Help
Real-time notificationsWatch for new orders and notify admins
Cache invalidationWatch for user updates and invalidate Redis cache
Search indexingWatch for new/updated posts and update Elasticsearch
Audit loggingWatch for all changes and log to an audit collection
Data syncWatch for changes and sync to another database
flowchart LR
Client[Client App] -->|INSERT user| Primary[Primary]
Primary -->|Write to| Oplog[(oplog)]
ChangeStream[Change Stream Listener] -->|Watches| Oplog
Oplog -->|Notify| ChangeStream
style Oplog fill:#f59e0b,color:#fff
style ChangeStream fill:#7c3aed,color:#fff

Change streams work by watching the oplog (operations log) — the same log that secondaries use for replication.


  • Change streams give you a live feed of database changes (inserts, updates, deletes)
  • Instead of polling (asking “anything new?” every second), MongoDB pushes changes to you
  • You can filter which changes you care about (e.g., only inserts)
  • Use resume tokens to recover from crashes without missing changes
  • Great for real-time features, cache invalidation, and data sync

Next: Performance & explain() →