Skip to main content
Mati Data, home

Writing

CDC that tells the truth: lag, drift and deletes

A change data capture pipeline can look healthy while it falls behind, loses a change or brings a deleted row back. The checks I build into my CDC pipelines, and what each one catches.

Most change data capture demos stop at the moment the first row lands in the warehouse. That is the easy part. The hard part is knowing, weeks later and with nobody watching, that the pipeline is still telling the truth.

CDC fails quietly. A stopped connector does not throw an error into anyone's inbox. A late event does not crash anything; it just overwrites a newer value with an older one. A delete that goes missing leaves a row in the warehouse that no longer exists anywhere else. Every service reports healthy, every dashboard is green, and the numbers are wrong.

This essay is about the checks I build into my CDC pipelines so that those failures announce themselves. Everything here is in public code you can run: CDC Starter is the smallest version, and the World Bank and Kiva pipelines are larger builds of the same pattern.

  1. PostgreSQL 17Explicit publication, a narrow CDC role, a heartbeat table
  2. Debezium 3.6.1Reads the write-ahead log through pgoutput
  3. RedpandaKafka-compatible event log
  4. ClickHouse 25.8Raw JSON, a typed view, ReplacingMergeTree by LSN
  5. dbtOne current-state macro, then staging and a mart
Quality, across every stage
10 dbt tests, one that fails if a deleted order returns; row counts and every order's state reconciled
Observability, across every stage
Heartbeat lag under a limit; connector and task state from Kafka Connect
The CDC Starter template: data moves top to bottom, and the checks that prove it arrived run across every stage.

Three ways a healthy pipeline lies

Each of these looks like success from the outside.

  • It stopped, and nothing seems to change. A quiet source table and a dead connector look the same in the warehouse: no new rows. If lag is measured from business tables, a quiet afternoon and an outage are indistinguishable.
  • A change arrived out of order, or twice. Message brokers deliver at least once, connectors restart and replay, and consumers fall behind and catch up. If the warehouse keeps whichever version arrived last, a redelivered old event can quietly win.
  • A delete vanished, or came back. Depending on how deletes are captured, a deleted row either never leaves the warehouse or returns with every column empty. Both pass a uniqueness test.

None of these raises an error. So the checks have to look for them on purpose.

Lag needs a clock that always ticks

Measuring lag as "time since the newest business row arrived" fails the first time the source goes quiet. The fix is a heartbeat: a one-row table that Debezium updates on a timer through its heartbeat query, which flows through the pipeline like any other change. Lag is then the age of the newest heartbeat in the warehouse, and it keeps moving whether or not anyone is writing orders. The heartbeat also keeps the replication slot advancing on a quiet database, so PostgreSQL can release the write-ahead log it no longer needs.

One detail matters more than it looks. When no heartbeat has arrived at all, the answer is not "zero lag". It is "no data", and it should fail. The Kiva monitor encodes this directly; simplified, it reads:

def compute_lag_seconds(now_utc, newest_source_update):
    """None means "no data has replicated yet", not "zero lag"."""
    if newest_source_update is None:
        return None
    return max((now_utc - newest_source_update).total_seconds(), 0)

An empty warehouse should never pass for a fresh one.

Order by log position, not by arrival

Every change in PostgreSQL has a position in the write-ahead log, the LSN. That is the only order that reflects what actually happened in the source. So I carry the LSN through as the row version and let the warehouse keep the highest one:

CREATE TABLE raw.orders
(
    order_id     UInt64,
    status       LowCardinality(String),
    _version     UInt64,   -- the PostgreSQL LSN of the change
    _is_deleted  UInt8
)
ENGINE = ReplacingMergeTree(_version, _is_deleted)
ORDER BY order_id;

Deduplicating on an ingestion timestamp looks identical in a demo and goes wrong the first time an old event is delivered again. Ordering by log position makes redelivery harmless: a late event carries an older version and loses.

Deletes are rows

By default, a delete arrives as a tombstone, a message with a key and no value, and many consumers simply skip it. I ask Debezium to rewrite deletes as ordinary rows with a flag instead, and to add the operation and the log position to every event:

"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.delete.tombstone.handling.mode": "rewrite",
"transforms.unwrap.add.fields": "op,lsn,source.ts_ms"

The warehouse table is told which rows are deletes, and reading with FINAL hides them. The catch is that two mistakes stay silent: forgetting FINAL returns duplicates, and forgetting the delete filter brings deleted rows back. So no model writes that read by hand. It lives in one dbt macro, and every staging model goes through it:

select *
from {{ source(source_name, table_name) }} final
where _is_deleted = 0

A singular dbt test then fails the build if a deleted order ever reaches staging. The test is small. The failure it catches is the kind that otherwise surfaces months later as a report nobody can reconcile.

Counts are not enough

Comparing row counts between the source and the warehouse catches lost and extra rows. It does not catch a change that has not arrived yet. Pause the connector and update an order's status: both sides still have the same number of rows, and the warehouse is wrong.

So the reconciliation compares state, not just size. In CDC Starter, every order's status is compared between PostgreSQL and ClickHouse after every run. For a small table that is a direct comparison. For a large one, the same idea scales by comparing aggregates or checksums per key range, then drilling into only the ranges that differ.

Ask the connector, not the throughput

Throughput graphs cannot tell an idle stream from a dead one. The Kafka Connect REST API can. I treat a connector as healthy only when the connector itself and every one of its tasks report RUNNING, and a connector with no tasks counts as down:

def connector_state_value(status_json: dict) -> int:
    """1 only if the connector itself and every task report RUNNING."""
    connector_ok = status_json.get("connector", {}).get("state") == "RUNNING"
    tasks = status_json.get("tasks", [])
    tasks_ok = len(tasks) > 0 and all(t.get("state") == "RUNNING" for t in tasks)
    return 1 if (connector_ok and tasks_ok) else 0

Layer the checks, then break them on purpose

No single check covers everything, which is why they are layered. CDC Starter runs four after every change: the connector and its task are running, row counts agree, every order has the same state on both sides, and heartbeat lag is under a limit.

The test I trust most is the one where I break the pipeline deliberately. Pause the connector and change nothing: the counts still agree, the connector check fails at once, and the heartbeat lag check fails once the limit passes. Now change a row while it is paused: the state comparison fails too. Resume the connector and every check passes again, with nothing lost, because the replication slot held every change in the meantime.

That drill is worth more than a green dashboard. It proves the alarms work. A check that has never failed in a test is a check you are trusting on faith.

What it costs

Very little. A heartbeat table and a narrow permission to update it. An explicit publication so only the intended tables are captured. A cap on how much write-ahead log a forgotten replication slot may hold, so a stopped consumer cannot fill the disk. One macro for reading current state, one test for deletes, and a short script that compares the two sides and asks the connector how it is.

Against that, the cost of a quiet failure is a decision made on numbers that were wrong for weeks.

Try it

CDC Starter runs the whole pattern on a laptop with one command, make demo, and the pause drill is in its README. The World Bank pipeline adds Prometheus alert rules with runbooks around the same checks, and the Kiva pipeline adds a reconciliation monitor that alerts by email. If you are building one yourself, the CDC pipeline checklist walks through each of these decisions.

References

  1. Debezium connector for PostgreSQL
  2. Debezium: new record state extraction (flattening events and handling deletes)
  3. ClickHouse: ReplacingMergeTree
  4. Kafka Connect REST API
  5. PostgreSQL: logical replication
  6. PostgreSQL: replication settings, including max_slot_wal_keep_size

See it in the code