Chapter 15 · Incremental Loading, CDC, Watermarks, High-Water Marks, and Idempotency
Idempotent Merge/Upsert Patterns, Batch Identity, Deduplication, and Retry Safety
Make retries safe with stable event and batch identity, deduplication, sequence-guarded upserts, tombstones, and atomic watermark advancement; distinguish idempotent processing from transport-level exactly-once claims.
Learning outcomes
After a worker crash, AtlasMart's change batch is delivered again. Two events are duplicated even within the successful delivery list. The pipeline must make a distinction: redelivery is normal transport behavior; applying a business effect twice is a target correctness bug.
Define idempotency as repeated application converging to the same target state, not as a synonym for deduplication.
Use stable event identity and source sequence to detect redelivery and conflicting duplicates.
Implement sequence-guarded current-state upserts and explicit delete tombstones in the local harness.
Keep batch identity, target mutations, audit state, and watermark transitions observable for retry/replay diagnosis.
Reject non-idempotent additive updates when the source event represents replacement state rather than a business delta.
Chapter 15 begins from the accepted Chapter 14 presentation state: eight current paid order-line facts, five paid orders, ten units, 690 USD paid GMV, 425 USD cost-at-sale, and 265 USD gross profit. The committed ERP change-stream cursor starts at source sequence 200. This chapter deliberately changes the current analytical state through six governed CDC events (sequences 201–206); after successful commit the target is nine paid lines, seven orders, eleven units, 740 USD GMV, 450 USD cost, and 290 USD gross profit. The change is explicit and reconciled rather than silently replacing earlier controls.
The mandatory lab is synthetic, local, and free. It uses
Python 3 standard library and its bundled
sqlite3 module. Generation-time validation ran
with Python 3.13.5 and SQLite 3.46.1; learners should record
their own versions. The source sequence is a simulated ordered
commit position, not a claim that every source database
exposes an identical integer cursor. The lab demonstrates
cursor, retry, ordering, tombstone, and schema-gate mechanics;
it does not reproduce a production write-ahead log, broker,
connector, distributed transaction, cloud service, or
transport-level exactly-once guarantee.
1. Idempotency is a state property
An operation is idempotent for a stated input identity and target contract when repeating it produces the same intended target state as applying it once. Deduplication is one technique: recognize the same event and skip its side effect. A sequence-guarded upsert is another: an older/repeated state cannot overwrite a newer one.
Batch identity answers a different question: which processing
attempt/range produced this state? The lab records
B20260921-CDC-A for the crashed attempt, B for the
successful restart, C for redelivery/reprocessing, and D for the
blocked schema event. Because A rolls back, its transactional
audit row also rolls back in this simplified harness; a
production control plane may retain failed-run telemetry
separately.
2. Event identity must detect conflicts, not just duplicates
row = conn.execute( "SELECT payload_sha256 FROM processed_event WHERE event_id=?", (e["event_id"],)).fetchone()if row: if row[0] != payload_hash(e): raise ValueError("conflicting duplicate event_id") return "duplicate"
If e202 arrives twice with identical payload, the second
delivery is harmless. If an upstream system reuses
event_id=e202 for different content, silently
ignoring it would hide corruption; the lab fails closed.
3. Sequence-guarded upsert prevents stale overwrite
INSERT INTO fact_sales_current( order_id,line_no,event_ts,customer_id,product_id,quantity, extended_amount,extended_cost,status,source_seq) VALUES (?,?,?,?,?,?,?,?,?,?)ON CONFLICT(order_id,line_no) DO UPDATE SET event_ts=excluded.event_ts, customer_id=excluded.customer_id, product_id=excluded.product_id, quantity=excluded.quantity, extended_amount=excluded.extended_amount, extended_cost=excluded.extended_cost, status=excluded.status, source_seq=excluded.source_seqWHERE excluded.source_seq > fact_sales_current.source_seq;
This exact syntax is SQLite-specific; MERGE/UPSERT syntax, multiple-match behavior, constraints, isolation, and retry semantics differ across engines. The vendor-neutral contract is “replace current state only with a strictly newer accepted source version for this business key.”
4. Controlled failure: increment a replacement event
-- WRONG if the CDC after-image says quantity is the complete replacement value.UPDATE fact_sales_currentSET quantity = quantity + :after_quantityWHERE order_id=:order_id AND line_no=:line_no;
On redelivery, the quantity increments again. The repair depends on source semantics: if the event is an after-image, assign the replacement value under a newer sequence; if the event is a true business delta, give the delta a stable identity and ensure it is applied at most once. Do not infer delta-vs-replacement from SQL convenience.
5. Deletes need idempotency too
e204 removes O1005/2 from current state and writes one tombstone keyed by delete sequence/event identity. Re-delivering e204 must not fail because the current row is already absent, nor create multiple logical deletes. The local current-state row count drops, while the tombstone proves the deletion was observed.
if current is not None and e["source_seq"] > current[0]: conn.execute( "DELETE FROM fact_sales_current WHERE order_id=? AND line_no=?", (oid, line) )conn.execute( "INSERT OR IGNORE INTO sales_tombstone VALUES (?,?,?,?)", (oid, line, e["source_seq"], e["event_id"]))
6. Retry acceptance evidence
The successful batch receives eight deliveries representing six
distinct logical events because e201 and e202 are redelivered.
Acceptance output reports applied: 6 and
duplicates: 2. A complete rerun after watermark 206
finds no event above the committed cursor and leaves the final
SHA-256 unchanged.
That proves convergence for the fixture, not universal “exactly once.” Production tests should include duplicate delivery across process restarts, stale/out-of-order events, conflicting identities, target deadlocks/timeouts, delete replay, and partial downstream publication.
Knowledge check
Check your understanding
- Is deduplication the same as idempotency?
- Why hash a processed event payload?
- What protects current rows from stale events?
- Why is quantity += after_quantity unsafe for an after-image?
- What does the Chapter 15 lab claim instead of exactly-once transport?
Review the answers
1. No. Deduplication is one mechanism; idempotency is the repeated-operation state property.
2. To detect reuse of the same event identity for conflicting content.
3. The source_seq guard only accepts a strictly newer source version.
4. Redelivery repeats the business effect even though the source state did not change again.
5. Deterministic target convergence under its stated at-least-once duplicate/retry fixture.
Summary and next step
This lesson established the mechanism and production boundaries for Idempotent Merge/Upsert Patterns, Batch Identity, Deduplication, and Retry Safety while preserving AtlasMart’s declared grain, governed metrics, history, and reconciliation evidence. Continue to Prove an Incremental Pipeline Handles Restart, Duplicate Delivery, Late Data, Deletes, and Reprocessing with those contracts unchanged unless an explicit, tested migration says otherwise.
Authoritative references
- Python documentation — sqlite3Local DB-API harness used to make commit/rollback and deterministic retry behavior observable.
- SQLite — TransactionTransaction boundaries used by the local acceptance lab; production engines have their own semantics.
-
SQLite — UPSERTExact local
ON CONFLICT ... DO UPDATEsyntax used for sequence-guarded convergence. - SQLite — DELETELocal current-state delete mechanics; the lab separately preserves a tombstone as CDC evidence.
- Debezium documentationOptional later-course reference for real CDC connector/envelope semantics. Debezium is not required by this chapter or local lab.
- PostgreSQL — Logical DecodingExample of a real database change-stream mechanism; PostgreSQL-specific details are not generalized to all engines.
- Kimball Group — Dimensional Modeling TechniquesDimensional grain/history semantics that incremental loading must preserve.