Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Incremental extraction: advance a compound watermark without losing ties

Last updated: 6 Oct 20265 min read
tutorial
IntermediateBy AITrove Editorial

A compound watermark combines update time and a stable tie-breaker so a paged extractor can resume without skipping records that share a timestamp.

Why one timestamp is insufficient

A source updates several receipts at the same millisecond. If a job stops after page one and resumes with updated_at greater than its last timestamp, every remaining receipt with the same timestamp is skipped. Order by (updated_at, receipt_id) and persist both values. Query lexicographically after the saved pair. Use stable source ordering and treat the checkpoint as part of committed output, not as an optimistic progress marker.

Define the consistency limit

A timestamp poll can still miss backdated changes, clock errors or deletes that leave no row. If those events matter, use a source change log or a bounded overlap window plus idempotent target upserts. The overlap trades repeated reads for a chance to catch late writes; it is not a proof against arbitrarily late records. CDC] carries a different change contract.

Commit output before progress

Write the selected page into a staging or target transaction and move the checkpoint only after the page is durable. On a crash before commit, reread the page; on a crash after commit, resume from the saved pair. Therefore the target write must tolerate retries. Idempotent loading] closes that loop.

Test the tie and crash

Create 53 source records with one identical update timestamp, page at 17 records, and interrupt after the second page. On restart, assert all 53 IDs appear once in the target. Then insert a backdated update and confirm the documented overlap or CDC policy handles it, or label that scenario unsupported.

Implementation

sql
SELECT receipt_id, updated_at, amount_cents
FROM receipt_source
WHERE updated_at > :last_time
   OR (updated_at = :last_time AND receipt_id > :last_id)
ORDER BY updated_at, receipt_id
LIMIT :page_size;

Performance and operating cost

With an index on (updated_at, receipt_id), each page can seek in O(log N) and scan O(P) rows for page size P. An overlap window increases repeated reads; a missing index turns polling into repeated table scans.

Common Mistakes

  • Do not advance on timestamp alone when records share it.
  • Do not save a checkpoint before its output is durable.
  • Do not assume timestamp polling observes physical deletes.

Read next

Continue the workflow: DAG intervals and idempotent task outputs.

ai-data
data-engineering
Storage details