Prompt 18 · Lesson 04 · Resumability and recovery
Resume Tokens, startAfter/resumeAfter, Invalidation, Oplog Windows, and Recovery
Resumability has a history horizon. A durable token is useful only while the deployment can still locate its event in retained oplog history.
Learning objectives
Persist and restore resume tokens without treating token internals as an application schema.
Use resumeAfter for normal continuation and startAfter after an invalidate event.
Explain why resumption depends on the token/timestamp still being locatable in retained oplog history.
Measure the current oplog window and compare it with consumer checkpoint age.
Recover from lost history through reconciliation plus a fresh stream checkpoint rather than skipping silently.
This lesson pins MongoDB Community Server
8.3.8 with
mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, mongosh 2.10.0, and PyMongo 4.17.0.
The mandatory lab uses a disposable
single-member replica set on loopback port
27158; change streams require a replica set or sharded cluster,
so a standalone is intentionally not used. A single-member
replica set is enough to learn event/resume mechanics but does
not demonstrate high availability, multi-node majority
durability, or failover. Authentication and TLS are disabled
only for this isolated lab. Feature Compatibility Version (FCV)
is inspected. FCV is observed and never changed. Default
read/write concern and primary read preference apply unless a
command says otherwise. The replica set uses a deliberately
modest 128 MiB oplog to make the retained-history concept
visible, but the mandatory lab does not generate enough data to
force rollover. A deterministic simulator demonstrates the
history-window decision without filling disks. Atlas, Search,
Vector Search, KMS, and Enterprise Advanced are not mandatory.
Product commands were not executed in this generation
environment because Docker, mongod, mongosh, and PyMongo are
unavailable here; runtime timings and token values must be
measured on the learner machine rather than copied as invented
output.
docker rm -f atlasmart-ch18-l4 2>/dev/null || truedocker volume rm atlasmart-ch18-l4-data 2>/dev/null || truedocker run -d --name atlasmart-ch18-l4 \ -p 127.0.0.1:27158:27017 \ -v atlasmart-ch18-l4-data:/data/db \ mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --replSet atlasmart-rs18-l4 --oplogSize 128 --bind_ip_alluntil mongosh "mongodb://127.0.0.1:27158/admin?directConnection=true" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27158/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-rs18-l4",members:[{_id:0,host:"atlasmart-ch18-l4:27017"}]})'until mongosh "mongodb://127.0.0.1:27158/admin?directConnection=true" --quiet --eval 'quit(db.hello().isWritablePrimary?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27158/admin?replicaSet=atlasmart-rs18-l4" --quiet --eval 'printjson(db.version());printjson(db.runCommand({getParameter:1,featureCompatibilityVersion:1}).featureCompatibilityVersion);printjson(rs.status().members.map(m=>({name:m.name,stateStr:m.stateStr})));'
1. A resume token is a checkpoint, not infinite retention
Each event token identifies a position in change-stream history.
resumeAfter opens a cursor after a normal event.
startAfter also starts after a token, but unlike
resumeAfter it can create a new stream after an
invalidate event. Both still depend on MongoDB
retaining enough oplog history to locate the referenced
operation. Store the entire BSON token; do not parse its private
_data representation into a business key.
from bson import json_utilfrom pymongo import MongoClientclient = MongoClient("mongodb://127.0.0.1:27158/?replicaSet=atlasmart-rs18-l4")db = client.atlasmartsource = db.orders_ch18_l4checkpoints = db.cdc_checkpoints_ch18_l4source.drop(); checkpoints.drop()source.insert_one({"_id":"o-18","status":"new"})with source.watch(max_await_time_ms=1000) as stream: source.update_one({"_id":"o-18"},{"$set":{"status":"paid"}}) event = stream.next() # Process the event first. Then persist the token. print(event["operationType"], event["documentKey"]) checkpoints.update_one( {"_id":"orders-consumer"}, {"$set":{"token":event["_id"],"clusterTime":event["clusterTime"]}}, upsert=True, )saved = checkpoints.find_one({"_id":"orders-consumer"})["token"]print(json_util.dumps(saved))with source.watch(resume_after=saved, max_await_time_ms=1000) as resumed: source.update_one({"_id":"o-18"},{"$set":{"status":"shipped"}}) print(resumed.next()["operationType"])client.close()
2. Invalidation is a different recovery branch
A collection-level stream is invalidated by operations such as
dropping or renaming the watched collection. The invalidate
event closes that cursor. Its token cannot be used with
resumeAfter; use startAfter if
starting a new stream after the invalidation is semantically
correct for the application.
const app=db.getSiblingDB("atlasmart");app.invalidate_ch18_l4.drop();app.createCollection("invalidate_ch18_l4");const s=app.invalidate_ch18_l4.watch();app.invalidate_ch18_l4.drop();const ev=s.next();printjson(ev);const token=ev._id;s.close();try { app.invalidate_ch18_l4.watch([], {resumeAfter:token}).next();} catch (e) { print(`resumeAfter invalidation is not allowed: ${e.codeName || e.message}`);}// startAfter is the API designed for beginning a new stream after invalidate.app.createCollection("invalidate_ch18_l4");const after=app.invalidate_ch18_l4.watch([], {startAfter:token});print("startAfter stream opened after invalidation");after.close();printjson(token);
The final startAfter stream is not forced to wait
for a new event in this static example; the important evidence
is the captured invalidate token and the rejected
resumeAfter path.
3. Measure the oplog window before trusting a checkpoint
const op=db.getSiblingDB("local").oplog.rs;const oldest=op.find().sort({$natural:1}).limit(1).next();const newest=op.find().sort({$natural:-1}).limit(1).next();printjson({ oldestTs:oldest.ts, newestTs:newest.ts, oldestWall:oldest.wall, newestWall:newest.wall, configuredBytes:op.stats().maxSize, usedBytes:op.stats().size});
from datetime import datetime, timedelta, timezonenow = datetime(2026, 9, 3, tzinfo=timezone.utc)oldest_retained = now - timedelta(hours=6)checkpoints = { "search-indexer": now - timedelta(minutes=2), "analytics": now - timedelta(hours=4), "abandoned-consumer": now - timedelta(hours=9),}for name, checkpoint in checkpoints.items(): can_resume = checkpoint >= oldest_retained print(name, {"checkpoint":checkpoint.isoformat(), "canResume":can_resume})
Expected deterministic result: the 2-minute and 4-hour checkpoints are inside a six-hour retained window; the 9-hour checkpoint is outside. This simulator is not MongoDB's token decoder—it makes the recovery decision explicit without manufacturing gigabytes of writes just to roll an oplog.
4. Wrong approach: “if resume fails, start from now”
Starting from now after a history-loss error silently creates a data gap. The repair is a reconciliation: compare/rebuild downstream state from the authoritative MongoDB snapshot, record a new synchronization boundary, then restart streaming from a supported fresh position. For search indexing that may mean reindexing all products; for a cache, rebuilding keys; for analytics, backfilling from durable source tables or another retained log.
def reconcile_products(source_collection, sink): source_ids = set() for doc in source_collection.find({}, {"_id":1,"name":1,"price":1}): source_ids.add(doc["_id"]) sink.upsert(doc["_id"], doc) for stale_id in sink.ids() - source_ids: sink.delete(stale_id) return {"sourceCount":len(source_ids), "sinkCount":len(sink.ids())}# After reconciliation succeeds, open a fresh stream and persist its new checkpoint.# Never claim that an expired old token was successfully resumed.
5. Production judgment
Monitor both consumer lag and checkpoint age versus current oplog window. The latter is the real recovery budget. Keep the same pipeline/options when resuming a token; MongoDB warns that changing them can make resumption unpredictable. Drivers automatically attempt a resume once for certain resumable errors, while mongosh does not; applications still need durable checkpoint storage, bounded retry/backoff, alerting, and a reconciliation path for non-resumable history loss.
Bridge. Lesson 5 uses these recovery mechanics to build idempotent downstream projections where a crash can replay an event without duplicating business effects.
Cleanup/reset
Everything in this lesson is disposable. Remove only the chapter-specific container and volume:
docker rm -f atlasmart-ch18-l4 2>/dev/null || truedocker volume rm atlasmart-ch18-l4-data 2>/dev/null || true
Check your understanding
- When is resumeAfter the normal choice?
- Why use startAfter after invalidate?
- What determines whether an old token can resume?
- What is unsafe about starting from now after history loss?
- What repairs history loss?
Review the answers
1. When continuing after a normal previously processed change event whose token is still in retained history.
2. resumeAfter cannot resume after an invalidate event; startAfter is designed to create a new stream after it.
3. The required operation must still be locatable in retained oplog history.
4. It silently skips changes between the old checkpoint and the new start point.
5. Reconcile/rebuild downstream state from an authoritative source, establish a fresh boundary, then resume streaming.
Authoritative references
- MongoDB Change Streams — deployment requirements, majority-committed notification, scopes, sharded behavior, resume tokens, and pre/post images.
- Change Stream Events — event fields, operation types, resume token, update/replace behavior, and expanded events.
- db.collection.watch() — pipeline stages, options, resumability, and mongosh versus driver behavior.
- db.watch() — database-scoped streams.
- Mongo.watch() — deployment-scoped streams.
- update Event — updateDescription, documentKey, full document and pre-image behavior.
- delete Event — delete event and pre-image behavior.
- invalidate Event — stream invalidation and startAfter boundary.
- Change Streams Production Recommendations — sharded total ordering and latency considerations.
- Privilege Actions — changeStream/find authorization requirements.
- Replica Set Oplog — retained history and oplog window.
- PyMongo Driver — official Python driver baseline and change-stream cursor APIs.
- PyMongo 4.17 Release Notes — current driver line used by the course.
- MongoDB 8.3 Release Notes — current stable minor and patch status.
- MongoDB 8.3 Compatibility Changes — current expanded-event field behavior inherited from 8.2.x.
- mongosh Release Notes — mongosh 2.10.0 baseline.