Prompt 18 · Lesson 05 · Idempotent integration
Build Idempotent Consumers for Search Indexing, Cache Invalidation, Outbox Alternatives, and Analytics
Real consumers crash between work and checkpointing. Make repeated events harmless, preserve business intent explicitly, and reconcile downstream state.
Learning objectives
Explain why at-least-once processing windows arise when side effects and checkpoint persistence are not one atomic action.
Build idempotent search/cache/analytics projections keyed by source identity or event identity.
Distinguish low-level CDC from an application-owned outbox/domain-event pattern.
Handle schema evolution, poison events, retries, and dead-letter/reconciliation paths without blocking forever.
Validate downstream state against MongoDB so a resumable stream is not mistaken for self-proving correctness.
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
27159; 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 lab simulates “search index”, cache,
and analytics projections with ordinary MongoDB collections so
no paid or external search/cache platform is required. Chapter
21 later covers MongoDB Search/Vector Search themselves. 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-l5 2>/dev/null || truedocker volume rm atlasmart-ch18-l5-data 2>/dev/null || truedocker run -d --name atlasmart-ch18-l5 \ -p 127.0.0.1:27159:27017 \ -v atlasmart-ch18-l5-data:/data/db \ mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --replSet atlasmart-rs18-l5 --oplogSize 128 --bind_ip_alluntil mongosh "mongodb://127.0.0.1:27159/admin?directConnection=true" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27159/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-rs18-l5",members:[{_id:0,host:"atlasmart-ch18-l5:27017"}]})'until mongosh "mongodb://127.0.0.1:27159/admin?directConnection=true" --quiet --eval 'quit(db.hello().isWritablePrimary?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27159/admin?replicaSet=atlasmart-rs18-l5" --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. Exactly-once side effects are not a change-stream guarantee
Suppose the search indexer receives event E, updates the search document, then crashes before writing E's resume token. After restart it resumes from the previous token and receives E again. That duplicate is expected in a safe at-least-once processing design. The consumer must make repeating E harmless or detect that E already affected the sink.
| Downstream | Useful idempotency key/pattern | Reconciliation |
|---|---|---|
| Search projection |
Upsert by source document _id; delete by
same key.
|
Compare source IDs/versions with index documents and rebuild drift. |
| Cache invalidation | Delete/invalidate is naturally repeatable; version cache fills. | Sample cache vs source or allow authoritative refill. |
| Analytics fact | Unique source-event ID / business-event ID prevents duplicate append. | Recompute counts/sums or compare event ledger. |
| Webhook/email | Persist delivery intent/idempotency key before external call. | Query provider/delivery ledger; do not rely on resume token alone. |
2. Reproduce the crash window and remove duplicate effects
from bson import json_utilfrom pymongo import MongoClientclient = MongoClient("mongodb://127.0.0.1:27159/?replicaSet=atlasmart-rs18-l5")db = client.atlasmartsource = db.products_ch18_l5search = db.search_projection_ch18_l5seen = db.consumer_seen_ch18_l5checkpoints = db.consumer_checkpoints_ch18_l5for c in (source, search, seen, checkpoints): c.drop()source.insert_one({"_id":"p-5","name":"Keyboard","price":70,"schemaVersion":1})with source.watch(full_document="updateLookup", max_await_time_ms=1000) as stream: source.update_one({"_id":"p-5"},{"$set":{"price":65}}) event = stream.next() event_id = json_util.dumps(event["_id"], json_options=json_util.CANONICAL_JSON_OPTIONS) def apply_idempotently(evt, eid): ledger = seen.find_one({"_id":eid}) if ledger and ledger.get("state") == "done": return "duplicate-suppressed" seen.update_one({"_id":eid},{"$set":{"state":"processing"}},upsert=True) doc = evt.get("fullDocument") # replace_one(upsert=True) converges on source _id, so replay is safe. search.replace_one({"_id":doc["_id"]}, doc, upsert=True) seen.update_one({"_id":eid},{"$set":{"state":"done"}}) return "applied" # Simulate a crash after the idempotent sink write but before marking done. seen.update_one({"_id":event_id},{"$set":{"state":"processing"}},upsert=True) doc = event.get("fullDocument") search.replace_one({"_id":doc["_id"]}, doc, upsert=True) # Replay sees an incomplete ledger and safely repeats the convergent upsert. print(apply_idempotently(event, event_id)) print(apply_idempotently(event, event_id)) # now suppressed because state=done checkpoints.update_one({"_id":"search"},{"$set":{"token":event["_id"]}},upsert=True) print(search.find_one({"_id":"p-5"}))client.close()
This local example demonstrates duplicate suppression, not a
universally atomic workflow. The seen, projection,
and checkpoint writes are separate operations here. If they must
share one MongoDB atomic boundary, a transaction can coordinate
them; if the side effect is external, use an external
idempotency key/delivery ledger and reconcile.
3. CDC versus an outbox/domain event
An outbox is appropriate when the application must publish
intent such as PaymentCaptured with a stable
business schema. The application writes the aggregate change and
outbox record in the same transaction or atomic aggregate
boundary. A change-stream consumer then watches the outbox
collection. This makes the event explicit; it does not magically
provide end-to-end exactly once.
const app=db.getSiblingDB("atlasmart");app.orders_ch18_l5.drop();app.outbox_ch18_l5.drop();app.orders_ch18_l5.insertOne({_id:"o-5",status:"pending",version:1});// Teaching shape only: in a real multi-document outbox, use a transaction on a replica set.app.orders_ch18_l5.updateOne({_id:"o-5",version:1},{$set:{status:"paid"},$inc:{version:1}});app.outbox_ch18_l5.insertOne({ _id:"evt-order-paid-o-5-v2", eventType:"OrderPaid", schemaVersion:1, aggregateId:"o-5", aggregateVersion:2, payload:{orderId:"o-5"}, createdAt:new Date()});printjson(app.outbox_ch18_l5.findOne());
The two writes above are intentionally shown as a shape, not as a correct multi-document atomic implementation. If losing the outbox record while committing the order is unacceptable, write both in a transaction (or remodel so the event marker lives in the same atomic document). The next-stage consumer still needs idempotency.
4. Poison events, schema evolution, and backpressure
A consumer that retries one schema-incompatible event forever can block its partition/stream and let the checkpoint age toward the oplog horizon. Validate event-envelope versions, distinguish transient infrastructure errors from permanent schema/data errors, quarantine poison events with the token and diagnostic context, alert, and keep a controlled reconciliation path. Never log entire sensitive pre-images merely for convenience.
SUPPORTED = {1, 2}def classify(event): doc = event.get("fullDocument") or {} version = doc.get("schemaVersion") if event["operationType"] == "delete": return "process-delete" if version not in SUPPORTED: return "quarantine-schema" return "process"
5. Reconciliation closes the correctness loop
Resume tokens prove stream position, not that every external projection is correct. Schedule reconciliation appropriate to business criticality: counts/hashes, sampled key/version comparisons, or full rebuilds. For search, compare source IDs and document versions; for analytics, recompute aggregates over a bounded window; for cache, rely on source-of-truth refill and versioned keys. Track duplicate rate, processing latency, checkpoint age, history-window margin, quarantine count, retry rate, and reconciliation drift.
source = { "p1": {"version":4,"price":10}, "p2": {"version":2,"price":20}, "p3": {"version":1,"price":30},}sink = { "p1": {"version":4,"price":10}, "p2": {"version":1,"price":19}, # stale "p4": {"version":1,"price":99}, # orphan}missing = sorted(set(source) - set(sink))orphan = sorted(set(sink) - set(source))stale = sorted(k for k in source.keys() & sink.keys() if source[k]["version"] != sink[k]["version"])print({"missing":missing,"orphan":orphan,"stale":stale})# Expected: missing=['p3'], orphan=['p4'], stale=['p2']
Production judgment. Change streams are excellent for near-real-time projections and integrations when consumers are resumable, idempotent, observable, and reconcilable. They do not replace a business-event contract, a backup, or a data-quality control. Size the oplog/recovery window for the longest realistic outage plus repair time; protect checkpoint stores; limit RBAC scope; keep external side effects idempotent; and rehearse recovery from both ordinary reconnects and history loss.
Bridge to Chapter 19. Time-series collections use an optimized storage model and do not support change streams. The next chapter therefore changes the data-model and retention mental model instead of assuming every MongoDB collection participates in the same CDC mechanism.
Cleanup/reset
Everything in this lesson is disposable. Remove only the chapter-specific container and volume:
docker rm -f atlasmart-ch18-l5 2>/dev/null || truedocker volume rm atlasmart-ch18-l5-data 2>/dev/null || true
Check your understanding
- Why can an event be delivered again after a crash?
- What makes search indexing naturally idempotent?
- When is an outbox preferable to inferring domain events from arbitrary writes?
- What should happen to a permanently schema-incompatible event?
- Why reconcile if resume tokens exist?
Review the answers
1. The side effect may complete before the consumer persists its resume checkpoint, so restart begins from an older token and replays the event.
2. Upserting a projection by stable source document identity/version makes repeating the same event converge to the same state.
3. When business intent and a stable event schema must be explicit and atomically tied to the aggregate change.
4. Quarantine it with diagnostic/token context, alert, and use an explicit repair/reconciliation path rather than retrying forever.
5. Tokens prove stream position, not correctness of external side effects or downstream state; reconciliation detects drift.
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.