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.

Intermediate–Advanced120–190 minutesCDC/resumability engineering labMongoDB 8.3.8 · mongosh 2.10.0 · PyMongo 4.17.0Last reviewed: September 2026

Learning objectives

01

Persist and restore resume tokens without treating token internals as an application schema.

02

Use resumeAfter for normal continuation and startAfter after an invalidate event.

03

Explain why resumption depends on the token/timestamp still being locatable in retained oplog history.

04

Measure the current oplog window and compare it with consumer checkpoint age.

05

Recover from lost history through reconciliation plus a fresh stream checkpoint rather than skipping silently.

Reproducible lab baseline

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.

start disposable replica set (l4)
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.

checkpoint every successfully processed event
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.

capture an invalidate token in an isolated collection
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

measure oldest/newest oplog timestamps
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});
deterministic checkpoint-vs-window decision
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.

reconciliation skeleton
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:

cleanup
docker rm -f atlasmart-ch18-l4 2>/dev/null || truedocker volume rm atlasmart-ch18-l4-data 2>/dev/null || true

Check your understanding

  1. When is resumeAfter the normal choice?
  2. Why use startAfter after invalidate?
  3. What determines whether an old token can resume?
  4. What is unsafe about starting from now after history loss?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.