Roadmap To Be A Data Engineer / Lesson 12

Lesson 12 CDC Fundamentals §1 Source Systems (+§4, §6) About 10 min read

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:

  1. Collapsed intermediates. A row that changed twice between photographs delivers only the second change.
  2. Invisible deletes. A deleted row is not in the table, so not in the result, so the warehouse keeps it forever.
  3. The commit-time gap. A transaction stamped updated_at = 02:58 that 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

  1. Which of our sources are polled on updated_at, and which read a log?
  2. Does anything in the source hard-delete rows? When was that table last reconciled against a direct count?
  3. Is the watermark set from the clock or from the data? What is the longest transaction on that table?
  4. Which reports depend on how long something stayed in a state?
  5. Can a bulk UPDATE on that table skip updated_at?
  6. 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)

  1. listings(id, state, price_cents, updated_at) + the polling extract, watermark in a file.
  2. Break it three ways on purpose: hard-delete 5% of rows; change rows twice between runs; open a transaction, UPDATE, wait past your extract, then COMMIT — confirm that row never returns.
  3. wal_level = logical; CREATE PUBLICATION; pg_recvlogical (or Debezium). Watch your own edits appear in commit order.
  4. Load the log into a bronze table; run the two queries from §6. Log-derived current state must equal SELECT * FROM listings exactly. Poll-derived will not.
  5. 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
Back to top