Roadmap To Be A Data Engineer / Lesson 22

Lesson 22 Incremental loads Fundamentals §5 Orchestration (+§1, §6) About 15 min read

The Mark the Tide Left

A watermark quietly drops rows from long transactions.

Incremental models & watermarks: how a job knows where it stopped.

An incremental model asks one question every run: where did I stop? Ask it with a business timestamp and the answer is a high-water mark — and a mark, once left, never comes down.

Model fct_listing_event
Datum max(event_ts)
Recording interval 15 min (was 60)
Observed loss 158,661 of 4,340,393 rows over 90 days (3.66%)

1The scene: 402 of twelve thousand

Tuesday 8 September, 09:20. The partner account manager forwards an email. Our largest wholesale supplier shipped a pallet of 12,000 second-hand garments on 27 August; the intake job logged 12,000 rows written; the shop’s admin panel lists 12,000 items live. Her onboarding dashboard, which reads the warehouse, says 402.

Nothing is broken. fct_listing_event has not failed a run in months. Freshness passes — the newest event is four minutes old. unique on event_id passes, so does not_null. No duplicates. Every dbt run since 1 August is green.

The model is ten lines long:

{{ config(materialized='incremental', unique_key='event_id') }}

select
    event_id, event_ts, listing_id, seller_id, event_type, price_cents
from {{ source('raw', 'listing_event') }}

{% if is_incremental() %}
  -- only what has happened since we last ran
  where event_ts > (select max(event_ts) from {{ this }})
{% endif %}

It has a unique key, so it cannot duplicate. It filters, so it does not re-scan history. The comment even says what it does. The comment is the bug.

2What a commit does to a maximum

A watermark is the model’s memory of where it stopped. This one is derived: recomputed every run as max(event_ts) over the rows already in the table. That makes it a maximum, and a maximum never goes backwards. Anything that arrives at or below the mark is not late. It is gone.

Rows arrive in the warehouse when the transaction that wrote them commits, not when they were stamped:

  • a shopper listing one item → a 1-row transaction, ~2 s;
  • a professional seller’s bulk importer → one row per item at ~1.4 s each;
  • the wholesale intake job → 1.05 s an item, so a pallet of 12,000 is one transaction running 3.50 hours.

Every row carries its own event_ts, spread across those hours — and all 12,000 become visible at the same instant, at commit.

Figure 1 — three transactions, one commit instant, one mark. All three commit at 00:40:00, become visible at 00:41:00, and are read by the 00:49 run, which sees max(event_ts) = 00:32:58:

Transaction Length Kept
wholesale intake, 12,000 rows 3 h 30 m 402 (96.65% dropped)
bulk lister, 600 rows 14 min 300 (50.00% dropped)
one shopper, 1 row 2 s 1 (nothing dropped)

11,598 of 12,000 rows were dropped by a where clause.

Replication lag does not cause this. A constant sixty-second lag delays every row equally and loses none. What clips rows is the spread between the first and last row of one transaction. Lag delays; variance in lag deletes.

3Completeness is a function of transaction length

August 2026, run every 15 minutes, 1,451,027 rows written by the shop (this table was produced in SQL, in DuckDB, against the generated event stream):

How long the transaction ran transactions rows written % that reached the model
under 1 minute 900,058 1,152,632 99.72
1 – 15 minutes 2,024 206,198 86.36
15 – 60 minutes 61 66,697 27.36
over 1 hour 3 25,500 3.89

Nine hundred thousand transactions are fine. Three are not.

The shape is not arbitrary. A transaction lands at some point between two runs, so the mark it meets is uniform over one run interval, which gives an expected survival of interval ÷ (2 × duration) once the transaction is longer than the interval. Measured across four orders of magnitude of transaction length (2 s → 12,600 s), the largest deviation from that closed form is 4.53 percentage points.

The left-hand end is structural, not lucky: a one-row transaction commits in about two seconds, well inside the sixty-second replication lag, so its timestamp cannot be below a watermark built from rows already visible one run earlier.

  • consumer (1-row) events: 2,662,587 over ninety days, 0 lost
  • professional sellers: 7.45% lost
  • the wholesale partner: 80.63% lost

4Fresher meant lossier

On 1 August the schedule changed from hourly to every fifteen minutes, to make the operations board fresher. Nobody thought of it as a data-retention setting. The same ninety days of source events, replayed at five cadences (same 4,340,393 rows every time):

Run every runs rows never arriving %
24 h 91 9,953 0.229
4 h 546 32,640 0.752
1 h 2,184 100,754 2.321
15 min 8,736 277,078 6.384
5 min 26,208 455,325 10.490

Monotone, and running the wrong way: the fresher the model, the more it loses. A shorter interval means a more up-to-date watermark, and a more up-to-date watermark sits further inside each long transaction. Hourly → 15 minutes multiplied the loss by ×2.75; five minutes would be ×4.52.

The honest baseline: the ticket did not create the bug. Over the sixty hourly days the model had already dropped 54,395 rows (1.88%, 906.6 a day); over the thirty days since, 104,266 (7.19%, 3,475.5 a day) — a rate ×3.82 higher. The wholesale intakes were already losing 47–70% of every pallet under the old schedule:

Pallet Items Length Cadence Kept Lost
2026-06-11 6,400 1.87 h 1 h 3,372 47.31%
2026-07-08 4,800 1.40 h 1 h 1,775 63.02%
2026-07-29 9,200 2.68 h 1 h 2,753 70.08%
2026-08-05 7,600 2.22 h 15 min 287 96.22%
2026-08-19 5,900 1.72 h 15 min 303 94.86%
2026-08-27 12,000 3.50 h 15 min 402 96.65%

5Right about the cash, wrong about the catalogue

Every money number in the company was exactly right.

  • August GMV €2,330,471.43 — identical to the cent in the source and in the model
  • 0 of 79,276 sale events lost
  • 47,833 listings created and never recorded

Not luck: arithmetic. A purchase is a one-row transaction — one buyer, one item, two seconds — and one-row transactions are the population this failure cannot touch. Sales, reservations and returns lost nothing. Bulk uploads lost 12.70% of new listings (47,833 of 376,518) and 7.71% of price changes (54,786 of 710,541).

By 30 August the warehouse was missing 70,434 listings that were live in the shop — 10.71% of the catalogue, against 24,248 (5.60%) on the day the schedule changed. Those listings carry €1,405,600.17 of listed value. They are not lost sales: they are on the site and can be bought. They are invisible to everything that reads the warehouse — the repricing job that marks slow stock down, the reverse-ETL sync to the marketing tool, the seller scorecards, the buying forecast, and the account manager’s onboarding dashboard.

6Why nothing fired

Every dbt test passed, and had to. unique, not_null, accepted_values, relationships are all predicates over rows that are in the table. A row that was never inserted satisfies every one of them. Freshness passed too — the newest event is always minutes old, because short transactions keep arriving.

Lesson 06’s trailing 30-day min–max band on the model’s daily row count crossed 7 times in 60 tested days, against 3 crossings on the true source counts, and both series first crossed on 2026-07-26, six days before the schedule changed.

On the worst day of the month:

rows z vs its own trailing 30 days
what the shop wrote, 27 Aug 56,241 +1.25
what the model kept 45,191 −0.20

The missing 11,050 rows are 2.35 standard deviations of the daily series, and the day still reads as an ordinary Thursday. The biggest data-loss day in three months turned a record day into an unremarkable one.

The one check that sees it is not clever — count both sides:

with src as (
    select date_trunc('day', event_ts) as d, count(*) as n
    from raw.listing_event group by 1),
mdl as (
    select date_trunc('day', event_ts) as d, count(*) as n
    from analytics.fct_listing_event group by 1)
select src.d, src.n as source_rows, coalesce(mdl.n, 0) as model_rows,
       src.n - coalesce(mdl.n, 0) as missing
from src left join mdl using (d)
where src.n <> coalesce(mdl.n, 0)
order by src.d

It returns a row for all ninety days — 60 of them before the schedule ever changed. Median shortfall 591 rows a day under the hourly schedule, 3,020 a day after, worst day 11,050. Its expected result is exactly zero, so it needs no threshold, no band, no tuning.

It was not skipped because it was expensive. The table is 3.26 GiB a year (17,520,000 rows); scanning two timestamp columns costs $0.0008 a run, $0.29 a year; scanning the whole table to rebuild and compare costs $0.0199. The model is incremental because a full rebuild takes longer than the gap between two runs, not because reading is expensive. It was skipped because nobody believed a copy could differ from its source.

7Six fixes, measured

August, 2,880 runs, 1,451,027 rows:

Configuration rows lost % write amplification
max(event_ts), append — today 104,266 7.19 1.00×
max(event_ts), merge on unique_key 104,266 7.19 1.00×
max(event_ts) − 30 min lookback, merge 22,386 1.54 2.90×
max(event_ts) − 4 h lookback, merge 0 0.00 16.81×
dbt microbatch, batch_size hour, lookback 1 (the default) 12,002 0.83 5.60×
dbt microbatch, batch_size day, lookback 1 0 0.00 134.26×
max(_ingest_seq), merge 0 0.00 1.00×

The unique key is not the fix. Merge and append lose exactly the same rows, because the filter runs before the merge does; the key only guarantees that whatever gets through is not written twice. Getting the key right is what made the model look healthy.

A lookback is a bet on the source’s longest transaction. Thirty minutes still loses 22,386 rows, 19,364 of them the wholesale partner’s; on the 27 August pallet it raises 402 items to 2,117. To lose nothing you need a lookback longer than the longest transaction plus one interval — here 225 minutes — a number that is a property of somebody else’s code and changes the day a supplier sends a bigger pallet. It is not free: on an append strategy a 30-minute lookback rewrites every row about three times (94.06% of rows exactly three times; 2,761,713 duplicate rows a month).

dbt’s own microbatch strategy is the same bet with better ergonomics. It keys on a config called event_time, “the column indicating at what time did the row occur”, and rebuilds the last lookback + 1 batches. The default lookback is 1. With batch_size: hour that still drops 12,002 rows in August at 5.60× the writing; the settings that lose nothing cost 17.53× (hour, lookback 4) or 134.26× (194,808,607 row-writes for 1,451,027 rows, at day granularity).

Watermark on something the writer controls, not something the row describes. A load-batch id, an ingestion sequence, or a table-format snapshot id (Lesson 17) is assigned when a row becomes visible, so it is monotone in visibility order — and a run always consumes a prefix of that order. The mark can never pass a row that has not landed. Zero rows lost, at any cadence, by construction rather than by tuning.

{{ config(materialized='incremental', unique_key='event_id',
          incremental_strategy='merge') }}

select * from {{ source('raw', 'listing_event') }}

{% if is_incremental() %}
  -- _ingest_seq is stamped by the loader when the row becomes visible,
  -- never by the application that wrote it
  where _ingest_seq > (select coalesce(max(_ingest_seq), 0) from {{ this }})
{% endif %}

Three conditions come with it:

  1. The sequence must be assigned by whatever makes rows visible. If a parallel loader hands out numbers before committing them, the same bug returns wearing a new column.
  2. The watermark is better stored as state in a run log than derived from the target, so a partial failure or a hand-run backfill cannot silently move it (Lesson 04).
  3. Reconcile anyway. The mark tells you where you stopped; the count tells you whether that was true.

8Ask the team

  1. What column does each incremental model watermark on, and who sets that column’s clock? “A timestamp from the source application” is the finding. You want “a number our loader assigns when the row lands.”
  2. What is the longest write transaction our source database produces, and who would tell us if it got longer? If nobody knows, no lookback in the project is defensible.
  3. Which incremental models have ever been compared with a full rebuild, and when? “Never” is the common and correct answer — which is why the daily count reconciliation is worth its 29 cents.
  4. When we changed the schedule, what else did we change? A cadence is a correctness setting in any timestamp-watermarked model, and belongs in the same review as a schema change.
  5. Which of our tests could fail if a row were missing? If the honest answer is none, we do not have completeness monitoring; we have shape monitoring.

9Hands-on (forty minutes, DuckDB or Postgres)

  1. Make a table of 200,000 events. Give 199,000 of them a commit_at one second after their event_ts. Give the last 1,000 a single shared commit_at and event timestamps spread over the previous three hours — that is your pallet.
  2. Write the loop by hand: repeatedly take max(event_ts) from your destination table, insert everything with commit_at <= now and event_ts > the mark, advance now by fifteen minutes. Ten lines of Python.
  3. Count both tables. Then count them grouped by event_ts::date. Then re-run the whole loop at one hour and at five minutes, and plot the three loss numbers.
  4. Change one thing: add a column seq numbered in commit_at order, watermark on that instead, and run all three cadences again. Every number should be zero.
  5. Add the reconciliation query to your project as a test with an expected value of zero, and check what it costs to run.

10Takeaway

An incremental model’s watermark must be monotone in the order rows become visible, not in the order the world made them. Business timestamps are not, because transactions take time; and the failure that follows is invisible to every test you own, because tests are predicates over rows that exist.

This is the thirteenth failure shape the series has collected, and the first with nothing to inspect. Events (06–10), the constant (11), the drift (12), the unwatched (13), the reversible (14), the unreproduced (15), the bundled (16), the transient (17), the extremal (18), the referential (19), the plural (20), the remote (21) — and now the absent. Every row in fct_listing_event is correct. The table is simply smaller than the truth, and no assertion over its contents can say so. The only detector is a second count from outside, which is precisely the computation incrementality exists to avoid.

Vocabulary

Term What it means here
watermark / high-water mark the value an incremental job stores to remember where it stopped. Derived (recomputed from the target) or stored (kept in a run log).
event time vs ingestion time when the world made the row vs when the row became visible to you. Only the second is monotone in arrival order.
commit-visibility spread the gap between the earliest and latest timestamp inside one transaction. This is what a timestamp watermark clips.
lookback deliberately re-reading a window below the mark. Requires a merge key to be safe, and its correct size is a property of a source you do not control.
incremental strategy append, merge, delete+insert, insert_overwrite, microbatch. Governs what happens to rows the filter selected — not which rows it selects.
write amplification rows written per row kept. A lookback buys completeness with it.
boundary reconciliation counting both sides of a hop and asserting the difference is zero. The only test that can see a missing row.

11How these numbers were made

Every figure comes out of one deterministic generator — no random seed, no wall clock — which simulates ninety days of listing events (4,340,393 rows) written in transactions of realistic length, plays the dbt run schedule over them, and records which rows the watermark filter would have selected. Running it twice gives the same numbers to the digit; the published artifact ships the generator as an appendix, and that copy was extracted from the rendered page, re-run in a clean directory and diffed key-by-key against the published figures (18 top-level keys, 814 scalar values, 0 differences).

  • The August event stream was written to CSV, loaded into DuckDB, and the reconciliation query, the GMV comparison, the per-seller and per-event-type splits and the transaction-length table were all run as SQL against it. The SQL and Python figures are asserted equal in the check script.
  • The generator asserts that the watermark is monotone, that loss rises with run frequency across all five cadences, that every cadence consumes the same 4,340,393 rows, that no sale event is ever clipped, and that consumer transactions lose nothing.
  • Prices, weekday factors and hour-of-day profiles are shaped to be realistic; the loss figures do not depend on them, only on transaction durations and the schedule.
  • The two dbt facts quoted — that microbatch keys on event_time and that lookback defaults to 1 — are read from dbt’s documentation, not from memory.

Previously: 04 idempotency and the logical run date · 12 how a poll loses a change · 13 counting both sides of a boundary · 15 an incremental model hiding its own blast radius · 17 snapshot ids as a version you can point at.

Back to top