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.

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

Learning objectives

01

Explain why at-least-once processing windows arise when side effects and checkpoint persistence are not one atomic action.

02

Build idempotent search/cache/analytics projections keyed by source identity or event identity.

03

Distinguish low-level CDC from an application-owned outbox/domain-event pattern.

04

Handle schema evolution, poison events, retries, and dead-letter/reconciliation paths without blocking forever.

05

Validate downstream state against MongoDB so a resumable stream is not mistaken for self-proving correctness.

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 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.

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

idempotent projection with checkpoint after successful sink update
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.

application-owned outbox record
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());
Atomicity warning

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.

version-aware consumer decision
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.

deterministic reconciliation check
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:

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

Check your understanding

  1. Why can an event be delivered again after a crash?
  2. What makes search indexing naturally idempotent?
  3. When is an outbox preferable to inferring domain events from arbitrary writes?
  4. What should happen to a permanently schema-incompatible event?
  5. 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

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.