Roadmap To Be A Data Engineer / Lesson 12
Between Two Polls
What a polling ingest silently misses, and why reading the log fixes it.
Synthetic marketplace: 37,020 listings, 2026-06-01 → 2026-08-29, nightly extract at 03:00. All figures produced by a fully deterministic generator (no RNG); the SQL in §6 is asserted against it with duckdb; the generator is shipped inside the artifact and was extracted back out of the rendered HTML and re-run (zero differing keys).
1The scene
Quarter end. Finance asks what the listed inventory is worth.
- Warehouse: €9,974,405.16 across 21,447 active listings.
- Production database, same instant: €8,495,189.43 across 19,920.
- Gap: €1,479,215.73 — the warehouse is 17.41% high.
Two proposed explanations, both rejected by arithmetic:
| Hypothesis | Check | Verdict |
|---|---|---|
| “The nightly job has been failing” | 90/90 extracts ran; daily net change in active listings inside its own trailing 30-day band on 26 of the last 29 days — the correct series scores 27 | rejected |
| “It’s a price or currency bug” | Quarter GMV €6,911,533.00 in the warehouse and €6,911,533.00 at source — equal to the cent across 15,041 sales | rejected |
Surviving explanation: the pipeline works exactly as designed, and the design cannot answer the question being asked of it.
2The mechanism
-- 03:00 every night
SELECT listing_id, state, price_cents, updated_at
FROM listings
WHERE updated_at > :watermark -- 03:00 yesterday
AND updated_at <= :now; -- 03:00 today
-- and afterwards
:watermark := :now;
It selects FROM listings — the table of current rows, not a record of changes.
Three consequences, none of them a bug in that SQL:
- Collapsed intermediates. A row that changed twice between photographs delivers only the second change.
- Invisible deletes. A deleted row is not in the table, so not in the result, so the warehouse keeps it forever.
- The commit-time gap. A transaction stamped
updated_at = 02:58that commits at 03:01 is invisible at 03:00 — and the watermark has already moved past 02:58, so it is never read again.
3Scale
Source committed 134,885 row versions. The extract delivered 67,841 (50.30%). Disjoint causes, asserted to reconcile:
| Fate | Versions | Share |
|---|---|---|
| Arrived | 67,841 | 50.30% |
| Overwritten before 03:00 | 41,061 | 30.44% |
| Committed after the read | 24,456 | 18.13% |
| Deleted at source | 1,527 | 1.13% |
A bug that heals only for well-behaved rows. 24,949 markdowns were applied; 2,288 (9.17%) reached the warehouse anyway because the item later sold and the sale re-touched the row. The other 22,661 (90.83%) never arrived — 22,283 of them on listings that are still active. The loss concentrates precisely on slow-moving inventory, which is exactly what a markdown job exists to reprice and what finance writes down at quarter end.
4The honesty check — the metric everyone watches could not break
GMV is not close, it is identical to the cent, and always would have been: a sale is a change that re-touches the row after the last thing that broke. Transactional metrics are structurally the ones a polling ingest gets right; state metrics are structurally the ones it gets wrong.
The count leak: 17.16 listings/day, crossing 1% on 2026-06-14 and 5% on 2026-07-15, reaching 7.67% by the end of the window.
Band test (Lesson 06’s instrument), daily net change in active listings:
| Series | Baseline band | Days inside |
|---|---|---|
| Warehouse | 139 … 264 | 26 / 29 |
| Source truth | 117 … 243 | 27 / 29 |
The corrupted series passes by the same margin as the clean one. Arithmetic, not luck: the leak adds ~17/day to a series whose ordinary day-to-day range is 125 wide.
A third shape of failure
- Lessons 06–10 detected events — a before and an after.
- Lesson 11 named the constant — invisible to every change detector because it never changes.
- This is a drift — each day’s increment is inside the band; the sum of ninety of them is not.
Threshold, band, anomaly and freshness tests are all detectors of change. A drift defeats them by being small per period; a constant defeats them by never occurring. Neither is caught by watching harder — both need an external reference: periodic reconciliation of warehouse counts against a direct count at the source.
5The sampling law
17,626 reservations; median lifetime 22.0 hours against a 24-hour poll. 64.17% resolve inside one polling interval; 4,957 began and ended between two consecutive 03:00 reads and left no trace at all. Of 15,511 listings ever reserved, the warehouse ever saw 11,650 (75.11%).
And the number still computes, plausibly and wrongly:
| Reservation fall-through rate | Seen | Fell through | Rate |
|---|---|---|---|
| From the change log (true) | 15,511 | 2,585 | 16.67% |
| From the polled warehouse | 11,650 | 2,201 | 18.89% |
Close enough to survive a sanity check, from a denominator a quarter too small, with errors that partly cancel. No error bar, no failing test.
Sharpen the question. “Is this number right?” is not answerable from inside the warehouse. “Could this number be right, given how the data got here?” is — and any metric that depends on how long something stayed in a state, or how many times something changed, cannot be computed from a daily poll, whatever value it returns.
6The fix — read the log the database already writes
Postgres WAL, MySQL binlog, SQL Server transaction log. Log-based CDC reads it instead of interrogating the table.
- Every version, not the latest one.
- Deletes are events (
op = 'D'), not absences. - Commit order is the only order (LSN); no window for a transaction to commit behind you.
- The source is not queried; latency drops from a day to seconds.
-- silver.listings — current state, rebuilt from the log
WITH latest AS (
SELECT *, ROW_NUMBER() OVER (PARTITION BY listing_id ORDER BY lsn DESC) AS rn
FROM bronze.listings_cdc
)
SELECT listing_id, state, price_cents, commit_ts
FROM latest
WHERE rn = 1
AND op <> 'D'; -- the delete wins the ranking, THEN removes the row
Order matters: a WHERE op <> 'D' applied before the ranking resurrects every deleted
listing — the polling bug, reintroduced in one line. Run against the 134,885 bronze rows this
returns 19,920 active worth €8,495,189.43 and 15,041 sales worth €6,911,533.00 — the source
figures, to the cent (asserted with duckdb).
-- validity intervals, straight out of the log (Lesson 09 for free)
SELECT listing_id, state, price_cents,
commit_ts AS valid_from,
LEAD(commit_ts) OVER (PARTITION BY listing_id ORDER BY lsn) AS valid_to
FROM bronze.listings_cdc
WHERE op <> 'D';
133,358 validity intervals, none of them inferred.
Three things that bite: the snapshot boundary (snapshot at LSN₀ then replay from LSN₀ —
overlap is fine if the apply is idempotent, a gap is silent loss); ordering is per key, not
global (partition by primary key); delivery is at-least-once (dedupe on (listing_id, lsn)).
7When polling is fine
| If the source table… | Polling is fine | You need the log |
|---|---|---|
| is append-only (events, orders, payments) | yes | only for latency |
| has mutable state you report on | no | yes |
| hard-deletes rows | no, unless you diff full snapshots | yes |
| changes faster than you poll | only if you want end state | yes, for anything about the path |
| is small enough to reload whole | yes — a full snapshot has none of these gaps | not for correctness |
has no trustworthy updated_at |
no | yes |
If a table fits in a nightly reload, reload it: no watermark, no watermark bugs, and
consecutive snapshots even recover deletes. Stuck on polling? Three changes remove most of the
damage: soft-delete instead of hard-delete; set the watermark from the data, not the
clock (max(updated_at) read, minus a lag longer than your longest transaction); and
reconcile against the source on a schedule.
8Ask your team
- Which of our sources are polled on
updated_at, and which read a log? - Does anything in the source hard-delete rows? When was that table last reconciled against a direct count?
- Is the watermark set from the clock or from the data? What is the longest transaction on that table?
- Which reports depend on how long something stayed in a state?
- Can a bulk
UPDATEon that table skipupdated_at? - What would tell us this was happening? (If the answer is “a monitor on the pipeline,” it would not.)
9Hands on (~90 min, one Postgres container)
listings(id, state, price_cents, updated_at)+ the polling extract, watermark in a file.- Break it three ways on purpose: hard-delete 5% of rows; change rows twice between runs; open a transaction,
UPDATE, wait past your extract, thenCOMMIT— confirm that row never returns. wal_level = logical;CREATE PUBLICATION;pg_recvlogical(or Debezium). Watch your own edits appear in commit order.- Load the log into a bronze table; run the two queries from §6. Log-derived current state must equal
SELECT * FROM listingsexactly. Poll-derived will not. - Compute “median time a listing spends in
reserved” from the log. Then try from the polled table, and write one sentence on why you cannot.
10Takeaway
Choose the ingest by the question you will be asked, not by the size of the table. If anyone will ever ask how long, how many times, or what changed, a poll cannot answer it — and it will return a number anyway.
11Vocabulary
| Term | Meaning |
|---|---|
| CDC | Change data capture — delivering a source’s changes rather than its current rows |
| query-based CDC | Polling a high-watermark column, usually updated_at. Blind to deletes, intermediate states and late commits |
| log-based CDC | Reading the WAL/binlog. Complete, ordered, low-latency; costs a replication slot and a connector |
| LSN | Log sequence number — the only trustworthy ordering key for applying changes |
| watermark | High-water mark of what has been collected. Safe from the data with a lag; unsafe from the clock |
| snapshot + stream | One full read, then the log from that exact LSN. The join is where CDC pipelines lose data |
| tombstone | The log entry recording a delete — the thing a poll can never produce |
| drift | A failure that adds a little each period. Passes every band/threshold test; only external reconciliation finds it |