MongoDB Change Streams in NestJS Keep Dropping Events. Here's the Fix
Change Streams look production-ready in a demo and fall apart on deploy day. Here's how to persist resume tokens, batch under load, and survive replica set elections in NestJS.
MongoDB Change Streams in NestJS Keep Dropping Events. Here's the Fix
MongoDB Change Streams are one of those features that look finished the moment they work. You wire up .watch(), log a change event to the console, and it feels like real-time architecture is basically solved. Then you deploy on a Friday, a pod restarts mid-shift, and you discover that every order update that happened during those four seconds of downtime just vanished. No error, no retry, nothing in the logs. It's just gone.
That gap between "works in a demo" and "survives production" is almost entirely about three things: resume tokens, backpressure, and reconnection. Get those three right and Change Streams are genuinely one of the best tools available for event-driven sync in a NestJS + MongoDB stack — better than polling, and lighter than standing up a full CDC pipeline with Debezium if your scale doesn't need it yet.
Resume tokens are the whole game
Every Change Stream event ships with an opaque _id — the resume token. It's MongoDB's bookmark. Pass it back in on your next watch() call via resumeAfter, and the stream picks up exactly where it left off instead of only listening from "now."
The failure mode is almost embarrassingly simple: if your app crashes or redeploys without saving that token somewhere durable, the next watch() call has no memory of it. It starts listening from the current moment. Whatever wrote to your orders collection during the restart window is invisible to your app forever — no error thrown, no exception to catch, because from Mongo's point of view nothing went wrong. You just weren't listening.
The fix is to persist the resume token to Redis (or any store that survives a restart) right after you've successfully processed each event — not before, because if processing fails you want to retry that same event on next boot, not skip past it:
@Injectable()
export class OrderChangeStreamService implements OnModuleInit, OnModuleDestroy {
private changeStream: ChangeStream;
constructor(
@InjectModel(Order.name) private orderModel: Model<OrderDocument>,
private redisService: RedisService,
) {}
async onModuleInit() {
await this.startListening();
}
private async startListening() {
const lastResumeToken = await this.redisService.get('mongo:resume:orders');
const options: ChangeStreamOptions = {
fullDocument: 'updateLookup',
...(lastResumeToken ? { resumeAfter: JSON.parse(lastResumeToken) } : {}),
};
this.changeStream = this.orderModel.watch(
[{ $match: { operationType: { $in: ['insert', 'update', 'replace'] } } }],
options,
);
this.changeStream.on('change', async (change: ChangeEvent<OrderDocument>) => {
try {
await this.handleOrderEvent(change);
await this.redisService.set('mongo:resume:orders', JSON.stringify(change._id));
} catch (err) {
console.error('Failed to process change event:', err);
}
});
this.changeStream.on('error', async (error) => {
console.warn('Change stream disconnected. Reconnecting with exponential backoff...', error);
this.reconnect();
});
}
onModuleDestroy() {
this.changeStream?.close();
}
}
One detail worth calling out: fullDocument: 'updateLookup' isn't optional if your downstream consumers need the full document on updates. Without it, an update event only gives you the changed fields, and you'd end up querying the database manually to reconstruct state — which defeats half the point of using a Change Stream in the first place.
At real throughput, one-event-at-a-time will hurt you
A Change Stream listener that processes events synchronously, one by one, works fine at low volume and falls over the moment you hit real traffic. A few thousand writes a second is enough to saturate the Node.js event loop and start generating unhandled promise rejections — not because the logic is wrong, but because you're doing too much synchronous work per tick.
The fix is boring but effective: buffer events in memory and flush in batches, either when the buffer hits a size threshold or a timeout fires, whichever comes first.
private eventBuffer: ChangeEvent[] = [];
private flushInterval = setInterval(() => this.flushBuffer(), 50);
private async handleIncomingEvent(change: ChangeEvent) {
this.eventBuffer.push(change);
if (this.eventBuffer.length >= 100) {
await this.flushBuffer();
}
}
private async flushBuffer() {
if (this.eventBuffer.length === 0) return;
const batch = [...this.eventBuffer];
this.eventBuffer = [];
await this.downstreamQueue.addBulk(batch.map(event => ({ name: 'sync-event', data: event })));
await this.redisService.set('mongo:resume:orders', JSON.stringify(batch[batch.length - 1]._id));
}
100 events or 50ms, whichever comes first, is a reasonable starting point for most order/inventory-style workloads — tune it against your own p99 write rate rather than treating it as a fixed constant. And notice the resume token save moved here too: you commit the token for the last event in the batch only after the whole batch has been successfully queued downstream, not per-event. That's the same "save after success" rule from the resume-token section, just applied at batch granularity.
Replica set elections will disconnect you — plan for it, don't just log it
MongoDB replica sets fail over. A primary steps down, an election happens, and for a window of a few seconds your Change Stream listener gets disconnected with something like MongoNetworkError or, worse, MongoError: resume token not found. This isn't a bug. It's a normal part of running a replica set, and a listener that just logs the error and dies is not production-ready.
Two things need to happen:
- Exponential backoff on reconnect. Don't hammer the primary the instant it errors — wait
Math.min(1000 * Math.pow(2, attempt), 30000)milliseconds before the next attempt, capping around 30 seconds. - Handle expired resume tokens explicitly. If the oplog has rolled over past your stored token — which happens if your app was down longer than the oplog retention window — MongoDB throws
InvalidResumeToken. At that pointresumeAftercan't help you. Log an alert, fall back tostartAtOperationTime, and treat it as a signal that you may need a full reconciliation sync to backfill whatever was missed.
That second case is the one teams tend to skip, because it only shows up when things have already gone wrong for a while — a long deploy freeze, an incident, a stuck pod. It's exactly the scenario where silently losing data hurts the most, so it's worth writing the fallback path before you need it, not after.
The checklist, if you're auditing an existing setup
| Operational requirement | Anti-pattern | Resilient solution |
|---|---|---|
| Downtime continuity | Reconnect with a fresh cursor (drops events) | Store the resumeAfter token in Redis after each event is processed |
| Event loop safety | Unbounded, synchronous per-event execution | In-memory batching with a timed flush |
| Failover handling | Unhandled stream crash on election | Exponential backoff reconnect listener |
| Data lookup overhead | Manually querying the DB for the full document | fullDocument: 'updateLookup' in stream options |
None of these are individually hard. What makes Change Streams painful in production is that all three failure modes are silent — no stack trace tells you a resume token expired three days ago, or that you've been dropping every fifth event under load. You find out from a support ticket about a missing order, not from your monitoring. Build the resume-token persistence, batching, and reconnect logic in from day one, and Change Streams stop being a liability and start being one of the more reliable pieces of an event-driven NestJS stack.
Accelerate your Backend & Database Modernization Roadmap
Need custom architecture auditing, automated OpenAPI contract generation, or zero-downtime microservice migration guidance for your engineering team?
Frequently Asked Questions
What happens if my NestJS app crashes while a Change Stream is open?
Nothing is queued for you. If you didn't persist the last resume token before the crash, MongoDB has no way to know where you left off, and a fresh watch() call only sees events from the moment it reconnects. Everything in between is gone.
Do I need Redis specifically to store resume tokens?
No, any durable store that survives a pod restart works — Redis is just fast and already in most stacks doing cache or queue duty. A small Mongo collection works too if you'd rather not add a dependency.
Why batch Change Stream events instead of processing them one at a time?
Because a naive listener processes events synchronously as they arrive, and at a few thousand writes a second that blocks the Node.js event loop and starts throwing unhandled promise rejections. Batching with a timed flush keeps throughput predictable.
Subscribe to RenovateAPI
Get weekly architectural guides, API refactoring strategies, and technical SEO updates delivered directly to your inbox.
Discussion (2)
Extremely helpful breakdown of the Strangler Fig pattern! We're currently refactoring a legacy Java monolith at work and the OpenAPI gateway routing tips saved us weeks of experimentation.
The schema JSON-LD and FAQ block structure really helps with indexing. Great technical detail on entity mentions too.
Suggested Related Articles
revalidateTag() Works Locally and Fails in Production: Fixing Next.js Cache Sync Across Multiple Nodes
Why revalidateTag() and revalidatePath() silently stop working once your Next.js App Router app runs on more than one server instance, and how to fix it.
Stripe Webhooks in NestJS Are Fine Until Production Hits: Fixing Race Conditions and Duplicate Events
How to stop Stripe webhooks from double-charging users, showing stale subscription status, or racing the frontend redirect in a NestJS and PostgreSQL app.
How to Modernize Legacy REST APIs: Architecture, Migration & Best Practices
A comprehensive guide on modernizing legacy REST APIs without breaking changes, covering strangler fig pattern, OpenAPI specs, and performance optimization.