Chapter 15 · Incremental Loading, CDC, Watermarks, High-Water Marks, and Idempotency

Timestamp/Sequence Watermarks, Overlap Windows, Clock Skew, Missing Rows, and Source Update Semantics

Design timestamp or sequence watermarks that cannot silently skip tied timestamps, clock-skewed updates, late business events, or ambiguous source update semantics; use overlap only with explicit deduplication and bounded assumptions.

Intermediate → Advanced125–150 minutesWatermark + overlap labPython 3 stdlib + sqlite3 · local/syntheticLast reviewed: September 2026

Learning outcomes

AtlasMart has six changes after cursor 200. If the engineer uses event_ts, the late O0999 sale appears to be four days old. If they use a source updated_at timestamp, sequence 205 carries 07:59:58Z even though sequence 204 was already stamped 08:04:00Z. A timestamp can be useful, but only when its generation and update semantics are known.

01

Distinguish event time, source update time, ingestion time, and ordered source sequence before choosing a watermark.

02

Explain how equal timestamps, clock skew, and non-monotonic source updates can create silent gaps.

03

Use composite cursors or overlap windows only when the source semantics make their guarantees explicit.

04

Show why overlap windows trade missed-event risk for deliberate rereads that require idempotent deduplication.

05

Choose a cursor from the source change contract, not from a conveniently indexed timestamp column.

Chapter 15 continuity contract

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.

Execution and guarantee boundary

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. Four clocks/positions that must not be collapsed

Field Question answered Safe as cursor when… AtlasMart use
event_ts When did the business event happen? Only if the source guarantees every later mutation has a later event time—rare for corrections/late data. Business attribution/history; not extraction cursor.
source_updated_at When did source logic stamp this row/change? Clock, precision, ties, update coverage, and monotonicity are governed. Counterexample: seq 205 is skewed behind seq 204.
ingested_at When did the pipeline observe it? For pipeline observability, not source completeness unless ingestion itself is authoritative. Lag/telemetry only.
source_seq Where is this change in the simulated ordered log? Source contract guarantees unique/ordered commit position and retention/recovery behavior. Primary Chapter 15 cursor.

A watermark is pipeline state over one chosen ordering domain. A high-water mark is the upper bound the batch intends to commit. The words are sometimes used differently in products; always document the exact field/domain.

2. The naive timestamp predicate misses a real change

Seq Entity/op Business event time Source updated_at Analytical effect
201 sale INSERT O1006/1 2026-09-21 08:00Z 08:01:00Z +60 GMV, +35 cost, +1 unit/order/line
202 sale UPDATE O1003/2 2026-09-19 15:00Z 08:02:00Z 50 → 55 GMV; cost remains 30
203 customer UPDATE C002 2026-09-21 08:03Z 08:03:00Z current segment Consumer → Growth
204 sale DELETE O1005/2 2026-09-20 07:10Z 08:04:00Z remove 100 GMV/60 cost/1 unit; preserve tombstone
205 late sale INSERT O0999/1 2026-09-17 13:00Z 07:59:58Z +80 GMV/+50 cost; source clock is behind seq 204
206 sale UPDATE O1002/1 2026-09-18 11:30Z 08:05:00Z 190 → 195 GMV; cost remains 125
naive_timestamp_incremental.sql
-- UNSAFE unless source_updated_at semantics guarantee monotonic completeness.SELECT *FROM source_changesWHERE source_updated_at > :last_timestampORDER BY source_updated_at;

If :last_timestamp = '2026-09-21T08:04:00Z', event e205 is invisible because its source timestamp is 07:59:58Z even though its ordered source sequence is 205. Generation-time acceptance output therefore prints timestamp-only cursor would miss: ['e205'].

3. Equal timestamps need a deterministic tie-breaker

Even with a trustworthy source clock, WHERE updated_at > last_ts can skip a row created with the same timestamp as the final row of the prior batch. One common relational technique is a composite cursor (updated_at, stable_unique_key) with lexicographic ordering:

composite_cursor.sql
SELECT *FROM source_tableWHERE updated_at > :last_ts   OR (updated_at = :last_ts AND source_pk > :last_pk)ORDER BY updated_at, source_pk;

This only works if source_pk is stable and the source's update behavior is compatible. If an existing row can be updated without changing updated_at, no tie-breaker repairs the fundamental source-contract failure.

4. Overlap windows deliberately reread

An overlap window moves the extraction lower bound backward, for example rereading changes since 07:59:00Z instead of only after 08:04:00Z. In the fixture that rereads e201–e206 and therefore captures clock-skewed e205. The benefit is tolerance to bounded delay/skew; the cost is duplicate delivery by design.

overlap_principle.py
cursor = "2026-09-21T08:04:00Z"overlap_start = "2026-09-21T07:59:00Z"  # explicit five-minute assumptioncandidate = [e for e in events if e["source_updated_at"] > overlap_start]# Every candidate must be deduplicated/idempotently merged downstream.

The five-minute number is not a recommendation. In production, derive an overlap from measured lateness/clock behavior plus policy, and monitor events that exceed it. If source updates can arrive arbitrarily late, a finite overlap cannot prove completeness.

5. Missing-row detection needs independent evidence

Watermarks tell the pipeline where it thinks it is. They do not independently prove that the source emitted every expected change. Use source counts/checksums where available, sequence-gap checks when sequences should be contiguous, periodic bounded snapshots, business control totals, or source/target reconciliation. In some logs sequence numbers are not contiguous because unrelated transactions share the log; gap rules must match the source.

Chapter 12 source contracts and Chapter 13 reconciliation therefore remain part of incremental design. A fast cursor with no completeness evidence can certify wrong data more quickly.

6. Lab: observe the skew and overlap tradeoff

watermark_demo.py
cursor = "2026-09-21T08:04:00Z"missed = [e["event_id"] for e in EVENTS          if e["source_seq"] > 204 and e["source_updated_at"] <= cursor]assert missed == ["e205"]overlap_cursor = "2026-09-21T07:59:00Z"reread = [e["event_id"] for e in EVENTS          if e["source_updated_at"] > overlap_cursor]assert "e205" in reread# Expected reread: e201..e206; downstream dedupe makes that safe.

What the evidence proves: this exact fixture defeats the naive timestamp cursor and the chosen overlap recovers it. What it does not prove: that five minutes covers production skew, or that every database exposes a total-order sequence suitable as the only cursor.

Knowledge check

Check your understanding

  1. Why is event_ts a poor generic extraction cursor?
  2. What is wrong with a bare updated_at > last_ts predicate?
  3. What does an overlap window require downstream?
  4. Can a finite overlap prove completeness under unbounded lateness?
  5. Why use source_seq in the local lab?
Review the answers

1. Corrections and late events can have old business times even when they are newly emitted.

2. Ties, clock skew, unstamped updates, and non-monotonic source behavior can cause gaps.

3. Stable event/business identity plus idempotent deduplication/merge because rereads are intentional.

4. No.

5. Its simulated contract provides a unique ordered change position, making the watermark mechanics observable.

Summary and next step

This lesson established the mechanism and production boundaries for Timestamp/Sequence Watermarks, Overlap Windows, Clock Skew, Missing Rows, and Source Update Semantics while preserving AtlasMart’s declared grain, governed metrics, history, and reconciliation evidence. Continue to Log-Based CDC: Inserts/Updates/Deletes, Before/After Images, Ordering, Transactions, and Schema Evolution with those contracts unchanged unless an explicit, tested migration says otherwise.

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.