When should a batch pipeline advance its stored extraction watermark to a new value?
answer
- only after the rows are really there
- late is cheap, early is unrecoverable
- not the wall clock on your host
- one transaction if the target allows
- a stuck value deserves an alert
basics
~20 sOnly after the target write for that batch has durably committed, and to a value derived from the data actually loaded or the window's upper bound — never to the current wall-clock time, and never before the load succeeds.
solid answer
~50 sThe watermark is a promise that everything below it is already in the target, so it may only move once that is true. Advance it **after** the load commits: if the job dies mid-write, the old value survives and the next run re-reads the window, which is safe because the write is keyed. Advance it **before** the write and a crash means the rows are gone forever, with no error anywhere. The value itself should come from the batch — the `hi` bound you read up to, or `max(updated_at)` of the rows you loaded — not from `now()`, which jumps past work that was in flight during the run. Best of all, write the state row in the same transaction as the load where the target supports it; where it does not, order it load-then-state and rely on the keyed write to make the resulting re-read harmless.
code
sql · 14 lines-- best case: the load and the bookmark move together or not at all
BEGIN;
MERGE INTO orders t
USING stg_orders s ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET status = s.status, updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (order_id, status, updated_at)
VALUES (s.order_id, s.status, s.updated_at);
UPDATE ingest_state
SET watermark = :window_hi, updated_at = current_timestamp
WHERE source_table = 'orders';
COMMIT;go deeper
Remember the order: load the rows first, then save the new position. Saving the position first means a crash loses those rows with no error.
Explain each crash point in the sequence and what the next run does about it, and why an idempotent keyed write is what makes load-then-state acceptable.
Show the operational side: per-table state, single-active-run locking, logged old and new values, an alert on a watermark that has stopped moving, and a safe rewind procedure.
Decide where pipeline state lives as a platform question — one queryable store, uniform semantics, auditable resets — so recovery does not depend on which team wrote the job.
## What the watermark actually claims A stored watermark (also called a bookmark, cursor or state value) is a durable claim: *every source row whose change value is below this is already in the target.* Every rule about when to advance it follows from taking that claim literally. If the value moves before the claim is true, the pipeline has lied to itself, and no future run will ever re-read that range — the rows are gone with no error raised. If the value moves later than it could have, the worst case is that some rows are re-read and re-applied, which a keyed write absorbs at the cost of a little compute. The asymmetry is total: late is cheap, early is unrecoverable. Design accordingly. ## Where the state lives The watermark must be as durable as the data it describes. In practice it lives in one of three places: - **A state table in the target warehouse**, keyed by pipeline and table — the most common choice, because it can often be updated in the same transaction as the load. - **The orchestrator's own metadata store**, if it offers durable per-task state. - **A small object in the same object store as the landed files**, written after the load. What it must not be: a variable in the job process, a file on ephemeral local disk, or something derived at runtime by querying the target (`SELECT max(updated_at) FROM target` looks clever and quietly breaks the moment the target holds rows from a backfill, a manual fix, or a different source). State should be per source table, not per pipeline. A single shared value forces the slowest table to gate the rest and makes a per-table replay impossible. ## The ordering rule The safe sequence is: read the window, land it to staging, apply the keyed write to the target, commit, **then** persist the new watermark. Every crash point in that sequence leaves the system in a state the next run repairs by re-reading: - Crash before the target write: nothing landed, watermark unchanged, next run reads the same window. - Crash after the target write but before the state update: rows landed, watermark unchanged, next run re-reads and the upsert makes it a no-op. - Crash after both: clean. The only genuinely bad ordering is state-first. That converts a crash into permanent silent data loss. ## Making it atomic where you can If the state table lives in the same database as the target, put the merge and the state update in one transaction and the intermediate case disappears entirely. Where the target is a warehouse that supports multi-statement transactions, this is the cleanest design and worth reaching for. Where the target cannot do that — object-store tables, systems without cross-table transactions, or a load performed by an external loader — accept at-least-once and lean on idempotency. That is the same reasoning that makes the lookback window affordable: if applying the same window twice is a no-op, you never need the two writes to be atomic. ## Choosing the value Three candidates, in descending order of safety: 1. **The window's upper bound**, when the run used explicit half-open `[lo, hi)` bounds and `hi` came from the source database's clock. This is unambiguous, reproducible, and does not depend on what the batch happened to contain. 2. **`max(change_column)` over the rows actually loaded**, computed from staging after the load. Correct, but it stalls on an empty batch — if no rows changed, the maximum is null and the watermark must simply stay put rather than being set to null or to zero. 3. **`now()` on the extract host.** Wrong. It moves the boundary past rows that were in flight during the run and past anything stamped by a slower clock, and it makes the run unreproducible. Whichever you choose, log the old value, the new value and the row count for the run. Watermark regressions and unexplained jumps are among the few pipeline bugs that are trivially diagnosable from logs — if the logs exist. ## Operating the state Treat the state store as a first-class operational surface. You want to be able to answer, in one query, "what is the watermark for every table and when did it last move?" A watermark that has not moved in three days is either a dead pipeline or a dead source, and both are worth an alert. Make manual reset possible but deliberate. Rewinding a watermark is the standard recovery for a missed window and must be safe by construction — with a keyed write it is; with an appending load it is a duplication incident. Guard it with a record of who reset what and to which value, because a rewind on a very large table is also a source-load event. Finally, never let two runs of the same pipeline hold the state concurrently. Overlapping runs — a slow run still going when the schedule fires again — will read overlapping windows and race on the update, with the loser's progress silently overwritten. A lock or a single-active-run policy is part of the state design, not separate from it.
- What should the watermark do when a run reads zero changed rows?Leave it where it is, or advance it to the window's upper bound if that bound came from the source clock. What it must never do is become null, zero, or the current time — the first two force a full re-read on the next run, the third skips whatever was in flight.
- Why is deriving the watermark from max(updated_at) in the target table a trap?Because the target is not only fed by this pipeline. A backfill, a manual correction, or a second source can plant a higher value, and the next run then skips the range between. Keep the pipeline's progress in its own state row so it means exactly one thing.
- Two runs of the same pipeline overlap because the first ran long. What breaks?Both read from the same starting watermark, so they read overlapping windows and race to write the state. Whichever finishes last wins, potentially rewinding or over-advancing progress. Enforce a single active run per table with a lock, and alert when a run overruns its interval.
saying these in an interview costs you the question
- Advancing the watermark before the target write commits
- Storing now() rather than a value derived from the data
- Keeping the bookmark in memory or on ephemeral disk
- Deriving the watermark by querying max() on the target
- Sharing one watermark across many source tables