Chapter 15 · Incremental Loading, CDC, Watermarks, High-Water Marks, and Idempotency
Log-Based CDC: Inserts/Updates/Deletes, Before/After Images, Ordering, Transactions, and Schema Evolution
Model log-based CDC as ordered source change evidence with inserts, updates, deletes, before/after images, transaction boundaries, and schema-version contracts rather than assuming delivery order equals database commit semantics.
Learning outcomes
Incremental polling asks “which rows look changed?” Log-based change data capture asks “which committed changes did the source record?” That can preserve deletes and ordered update evidence much better—but only if the consumer understands the event envelope, transaction ordering, schema evolution, and retention/recovery contract of the actual source/connector.
Define a CDC event envelope and distinguish insert, update, delete, before-image, and after-image semantics.
Separate source commit ordering from network/delivery ordering and explain why consumers need an explicit ordering field.
Explain transaction-boundary requirements when one source transaction emits multiple row changes.
Treat schema version as part of the stream contract and block unsupported breaking changes before cursor advancement.
Distinguish the local simulated log from real database WAL/binlog/logical-decoding or connector guarantees.
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. CDC is a stream of change evidence, not a magical table mirror
Change data capture (CDC) is a mechanism for obtaining inserts, updates, deletes, or transaction changes from source change evidence. Log-based CDC derives those changes from a database transaction/write log or an equivalent ordered change facility rather than comparing snapshots or polling a business timestamp column.
The Chapter 15 local fixture simulates an ordered log with
fields such as event_id, source_seq,
schema_version, entity,
op, before, and after.
Real products use different names, transaction metadata, offset
domains, retention policies, and delivery guarantees.
| Envelope field | Purpose | Important boundary |
|---|---|---|
| event_id | Stable delivery identity for dedupe/conflict detection. | Must not be regenerated on every redelivery if used as dedupe identity. |
| source_seq | Ordered source change position in this simulation. | Not equivalent to every vendor LSN/binlog offset. |
| op = I/U/D | Change kind. | Deletes may have only key/before image depending on source. |
| before | Prior values when available. | Availability can depend on source configuration and connector. |
| after | New values for inserts/updates. | A schema change can alter required fields. |
| schema_version | Contract version understood by consumer. | Unsupported breaking version blocks batch in the lab. |
| event_ts | Business-time semantics. | Not used as change ordering cursor. |
2. Before/after images support more than current-state upsert
For e202 the before image says O1003/2 had GMV 50 and the after image says 55. A current-state table only needs the after image, but before images can help compute deltas, audit corrections, validate expected prior state, or feed downstream change-aware consumers. They can also contain sensitive data that current rows no longer expose.
For delete e204, the current fact row is removed but a tombstone records the business key and source sequence. This prevents “the row vanished” from being mistaken for “the pipeline never saw it.” Whether production history physically deletes, soft-deletes, anonymizes, or retains a separate ledger is a governance decision tied to Chapters 11 and 22.
3. Source order and delivery order are different
The local delivery list intentionally contains duplicates and is
not trusted as authoritative order. The consumer sorts accepted
candidates by source_seq. In a real log, ordering
scope can be global, partition-local, transaction-local, or
connector-specific. Do not sort by network arrival time and call
that database order.
Transaction boundaries matter too. If one source transaction changes an order header and two lines, publishing/committing only part of that transaction can create a state the source database never exposed. Production CDC consumers may need transaction metadata or source-specific atomicity rules. The local fixture uses one logical row change per simulated source sequence, so it does not reproduce multi-row source transaction assembly.
4. Schema evolution is part of offset safety
After target convergence through sequence 206, the lab presents
e207 with schema_version=2 and a renamed
extended_amount_cents field. The consumer supports
only version 1. The safe response is to reject the candidate
transaction and leave both target hash and watermark at 206.
def validate_event(e): if e.get("schema_version") != 1: raise ValueError( f"unsupported schema_version={e.get('schema_version')} " f"at source_seq={e.get('source_seq')}" )# e207 is rejected before commit; watermark remains 206.
Adding an optional field may be compatible for one consumer; renaming an amount, changing units, changing null semantics, or changing key meaning can be breaking even when JSON still parses. Chapter 12's contract principle applies equally to streams.
5. Controlled failure: treat at-least-once delivery as exactly-once source change
A connector or broker may redeliver after retry/rebalance. If
every delivery executes INSERT or adds deltas
again, business effects multiply. Conversely, claiming “exactly
once” because the target currently has no duplicates confuses an
observed outcome with an end-to-end guarantee.
The Chapter 15 local guarantee is narrower and testable: duplicate event IDs are recognized; sequence-guarded current-state writes are idempotent; target mutations and watermark advancement share one SQLite transaction. That produces convergence under this fixture's at-least-once redelivery without asserting exactly-once transport.
6. Source-specific verification before production
Before implementing real log-based CDC, verify from the source/connector's current official documentation: offset meaning, ordering scope, transaction metadata, before/after image behavior, delete/tombstone representation, snapshot-to-stream handoff, retention, restart semantics, schema-change behavior, and security exposure. The optional PostgreSQL and Debezium references below illustrate real mechanisms, but their exact semantics must not be projected onto another database.
Knowledge check
Check your understanding
- What makes CDC different from polling updated_at?
- Does network arrival order prove source commit order?
- Why preserve a delete tombstone?
- What happens to e207 in the local lab?
- Does this lab reproduce a real WAL?
Review the answers
1. CDC consumes explicit source change evidence, often including deletes/order/transaction metadata that a timestamp scan may not provide.
2. No; use the source/connector ordering contract.
3. It records that a governed delete was observed and gives retry/audit logic stable evidence.
4. It is rejected as an unsupported schema version; target and watermark remain at the last committed state 206.
5. No. It is a deterministic simulation of consumer-side CDC mechanics.
Summary and next step
This lesson established the mechanism and production boundaries for Log-Based CDC: Inserts/Updates/Deletes, Before/After Images, Ordering, Transactions, and Schema Evolution while preserving AtlasMart’s declared grain, governed metrics, history, and reconciliation evidence. Continue to Idempotent Merge/Upsert Patterns, Batch Identity, Deduplication, and Retry Safety 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.