skip to content

How does a watermark-based incremental snapshot backfill a table while the CDC stream keeps running?

level: seniorimportance: nice to knowfreq 38%

answer

  1. use the log itself as a clock
  2. two markers bracket the risky interval
  3. the chunk is read between them
  4. keys changed inside the window are dropped
  5. the stream already carries the truth for those

basics

~20 s

The connector writes a marker into the log, reads one key-range chunk, writes a second marker, then drops from the chunk any key that changed between the markers — the live stream already carries the truth for those. Streaming never pauses.

solid answer

~50 s

Chunk the table by primary key and reconcile each chunk against the log instead of freezing the database. For a chunk, the connector writes a **low watermark** row into a dedicated signalling table — that write appears in the log as an ordinary event — then selects the chunk into memory, then writes a **high watermark** row. As it continues processing the log, any change event for a key that falls between the two watermarks means the in-memory copy of that key is stale, so that key is removed from the chunk. The surviving chunk rows are emitted as baseline events. The result: no long-running transaction, no locks, no paused stream, per-chunk resumability, and the ability to backfill a newly added table or re-seed one table without touching the rest. It is the algorithm behind incremental snapshots in modern CDC frameworks.

code

text · 13 lines
text
log stream (single ordered sequence)
  ... normal change events ...
  LOW  WATERMARK  chunk=41        <- connector writes to its watermark table
                                     (SELECT of pk 1000..1999 happens here)
  UPDATE orders id=1042 status=SHIPPED
  DELETE orders id=1777
  HIGH WATERMARK  chunk=41
  ... normal change events ...

in-memory chunk read: 1000..1999      (1000 rows)
keys changed in window: {1042, 1777}  -> removed from the chunk
emitted as baseline:    998 read events
1042 and 1777 reach the sink from the stream, with current values

go deeper

for a junior

Recognise the shape: the table is backfilled a chunk at a time while changes keep flowing, rather than the whole table being read in one frozen pass.

for a middle

Explain the two markers and what the interval between them means — the window in which the rows you just read could have gone stale — and why keys changed in that window are dropped from the chunk.

for a senior

Argue why the markers must be real writes in the log rather than clock readings, and name what the design buys operationally: no long read view, per-chunk resumability, and the ability to re-seed one table alone.

for a principal

Weigh the prerequisite honestly — a writable table inside a source you may not own — against the alternative of a locking or long-read-view snapshot, and decide whether losing a single-instant cross-table baseline is acceptable for your consumers.

## The problem it solves A classic snapshot freezes an instant: one long read view, one pinned log position, streaming starts afterwards. That works but is brittle at scale — the read view blocks version cleanup, the pinned position pins the log, and you cannot add a table to the capture set later without redoing the whole thing. What you actually want is to backfill history **while** the stream runs, in resumable pieces, with no lock and no long transaction. The obstacle is that a chunk read at 10:00 and emitted at 10:02 may already be stale — the stream may have carried an update to one of its rows in between — and emitting the stale row after the fresh one drives that row backwards. ## The watermark trick The insight is that the log itself can be used as a clock. The connector owns a small table in the source (often called a signal or watermark table) whose sole purpose is to be written to. A write to it produces a log event like any other, and that event's position in the stream is an unambiguous ordering marker relative to every other change. For each chunk the connector performs: 1. **Write low watermark.** Insert or update a row in the watermark table with a chunk identifier. This event will appear in the stream. 2. **Select the chunk.** `SELECT ... WHERE pk > :last ORDER BY pk LIMIT :size` — an ordinary short read, held in memory as a key-to-row map. 3. **Write high watermark.** A second write to the watermark table with the same identifier. 4. **Keep processing the log.** Events keep flowing to the sink normally. When the low watermark event is seen, the connector starts noting keys: any change event between the low and high watermark whose key is in the in-memory chunk causes that key to be **deleted from the chunk**, because the stream is about to deliver (or has delivered) a newer version. 5. **Emit on the high watermark.** When the high watermark event arrives, whatever survives in the map is emitted as baseline read events. Everything the chunk emits is therefore either untouched during the window — safe — or has been superseded by an event the stream carries. No row goes backwards. ## Why the ordering argument holds The watermark writes are ordinary transactions, so their positions in the log are total relative to other commits. The window between them is exactly the interval during which the in-memory copy could have become stale. Anything committed before the low watermark is already reflected in the chunk read; anything after the high watermark arrives at the sink after the chunk is emitted, in the normal stream order. The only ambiguous interval is the one the algorithm explicitly reconciles. ## What it buys operationally - **No pause and no long transaction.** Latency for live changes stays flat during a backfill of a multi-terabyte table, and version cleanup on the source is never blocked for more than a chunk. - **Resumable and pausable.** Progress is per chunk. A restart resumes at the next boundary; a backfill can be throttled or stopped during peak hours and continued overnight. - **Targeted.** A single table newly added to the capture set, or one table whose sink was corrupted, can be re-seeded without re-snapshotting anything else — the classic snapshot's all-or-nothing weakness. - **Bounded memory.** Only one chunk is resident. ## What it costs - **A writable table in the source.** The connector must be able to create and write a small table in the database being captured, which some data owners refuse outright. That refusal is the usual reason teams fall back to a classic snapshot. - **A usable chunking key.** Ordered chunking needs a primary key that can be ranged over. Composite keys work; random UUID keys still chunk correctly by ordering but give unhelpfully scattered I/O, and a table with no key at all is a problem. - **Extra log volume.** Two extra events per chunk, plus a small write amplification on the source. - **Longer overall convergence.** Backfilling gently over days is precisely the point, but the sink is incomplete until it finishes, and consumers must tolerate that. ## The failure mode to name If the chunk is emitted **without** the reconciliation step — for example, a home-grown backfill that reads chunks and writes them straight to the sink alongside a running stream — then a chunk row read before a concurrent update lands after that update at the sink, and the row silently reverts to its old value. The corruption is invisible until someone compares counts by value. The watermarks exist exactly to make that impossible.

  • Why must the watermarks be written to a real table rather than tracked in the connector's memory?
    Because their only job is to appear in the log, where their position is totally ordered against every other commit. An in-memory timestamp cannot be compared reliably with commit order — clocks drift, and events are not stamped with the connector's clock. Writing a row makes the source's own log assign the ordering, which is the guarantee the algorithm depends on.
  • What happens if a key in the chunk is deleted in the source during the reconciliation window?
    The delete event falls between the watermarks, so the key is removed from the in-memory chunk and never emitted as a baseline row. The stream delivers the delete itself, and the sink either removes the row or marks it deleted. Emitting the chunk row instead would resurrect a deleted record — the same reversal failure the watermarks exist to prevent.
  • When would you still choose a classic snapshot over this approach?
    When you cannot create a writable table in the source, which is common on databases owned by another team or a vendor. Also when the tables are small enough that a single consistent read finishes in minutes, or when you need a genuinely single-instant baseline across many tables for cross-table consistency at the sink — the incremental approach deliberately gives that up in exchange for never pausing.

It is like taking inventory of one shelf while the shop stays open: you note the time you started and finished that shelf, and for any item sold in between you trust the till receipt rather than your count.

saying these in an interview costs you the question

  • Thinks the stream must pause while a chunk is read
  • Uses wall-clock timestamps instead of log positions to order the window
  • Emits chunk rows without reconciling keys changed during the read
  • Assumes a table with no usable key can be chunked
  • Believes the technique removes the need for idempotent sink writes

context