- 4 checks after every run: connector state, row counts, the state of every order and heartbeat lag
- 10 dbt tests, one of which fails if a deleted order ever reaches staging
Context
My two public CDC pipelines, one for World Bank indicators and one for Kiva loans, carry a lot around the core pattern: an API client, an orchestrator, dashboards and alert rules. When someone asks how to start their own, that is too much to read before the first row moves. CDC Starter is the pattern with everything else taken out, and with the parts that usually go missing put in from the start: deletes, ordering by log position, lag measured against a heartbeat, and a check that the warehouse really matches the source.
It is built to be copied. make demo starts PostgreSQL, Debezium, Redpanda and ClickHouse, registers the connector, writes inserts, updates and a delete to the source, waits until ClickHouse shows exactly the same orders, builds and tests the dbt models, and runs the checks. It is the running example behind the CDC workshop and the step-by-step build guide on this site.
Constraints
- One command, and safe to run again: a second
make demopasses like the first. - Laptop-sized. Every image is pinned, every container has a memory cap, about five gigabytes in total, and ports bind to loopback on non-default numbers, so the stack cannot collide with a database already running on the machine.
- Nothing to install but Docker,
makeandcurl. dbt runs in its own pinned image, so nobody needs a local Python environment. - Small enough to read: one table, one connector, two dbt models, and checks written as a shell script with a meaningful exit code instead of a dashboard.
Architecture
- PostgreSQL 17Explicit publication, a narrow CDC role, a heartbeat table
- Debezium 3.6.1Reads the write-ahead log through pgoutput
- RedpandaKafka-compatible event log
- ClickHouse 25.8Raw JSON, a typed view, ReplacingMergeTree by LSN
- 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 source is an app.orders table in PostgreSQL 17 with logical replication on and a trigger that moves updated_at on every change. The publication lists only the tables meant for capture, and Debezium connects as its own role, with replication and read access and the one update its heartbeat needs.
Debezium 3.6.1 reads the write-ahead log through pgoutput and publishes flattened rows to Redpanda. Each event carries the operation, the log position (LSN) and a delete flag, and a delete arrives as a row of its own instead of a bare tombstone.
ClickHouse 25.8 reads each message through a Kafka engine table as one raw JSON string. A materialized view types it and writes raw.orders, a ReplacingMergeTree versioned by the LSN and told which rows are deletes, so the latest change in the source wins rather than the latest to arrive.
dbt reads the current state through one macro, FINAL plus the delete filter, into a staging view, then builds a mart of orders by status. Source freshness on the raw table catches a stopped pipeline while the mart still looks fine. make check then reconciles the source and the warehouse, measures lag against the heartbeat, and asks Kafka Connect whether the connector and its task are running.
Decisions and tradeoffs
- Order by log position, not by arrival. The version column is the PostgreSQL LSN, so a late or redelivered event cannot overwrite a newer one. Deduplicating on an ingestion timestamp looks the same in a demo and goes wrong in production.
- Deletes are rows. The connector rewrites a delete as a row with a delete flag, and
ReplacingMergeTreehides it underFINAL. Without this, a delete either vanishes or comes back as a row with empty columns. - The current-state read lives in one macro. Forgetting
FINALgives duplicates, and forgetting the delete filter brings deleted rows back. Neither raises an error, so no model writes them by hand, and a singular test fails the build if a deleted order ever reaches staging. - Raw JSON first. The Kafka engine reads a string and the view types it, so a malformed message is skipped instead of stalling the consumer. The cost is that nothing enforces a schema at the edge, which a schema registry would.
- An explicit publication and a narrow role. Debezium is told not to create its own publication, so adding an unrelated table to the database never starts capturing it by surprise.
- A heartbeat for lag. Debezium updates a heartbeat row on a timer, which measures lag even when the business tables are quiet and keeps the replication slot moving. No heartbeat at all counts as a failure, not as zero lag.
- The slot survives restarts. The connector keeps its slot when it stops, so a restart resumes from the recorded position. A cap on the write-ahead log a slot may hold stops a forgotten slot from filling the disk, and
make downremoves everything, the slot included.
Outcome
make demopasses from a clean start, aftermake downhas removed every container and volume, and passes again when run a second time.- 4 checks after every run, each printing PASS or FAIL: the connector and its task running according to the Kafka Connect API, row counts that agree, the same state for every order on both sides, and heartbeat lag under a limit.
- 10 dbt tests: keys, accepted values and not-null checks on the staging view and the mart, plus the singular test that fails if a deleted order reaches staging.
- The failure path is tested, not only the happy one. With the connector paused and nothing written, the row counts still agree: the connector check fails at once, and the heartbeat lag check fails once the limit passes. Changing the source while paused fails the reconciliation as well. After a resume every check passes again and nothing is lost, because the replication slot held the changes.
What I would change
- A make target that runs the pause drill by itself, so CI exercises the failure path and not only the happy one.
- The checks exported as Prometheus metrics with alert rules, as in the two larger projects, once the template is used for more than learning.
- A schema registry with Avro or Protobuf instead of plain JSON, once more than one team writes to the source.
- A second table and a join in the mart, to show how
FINALand the delete filter behave across tables. - An orchestrator for dbt, Airflow or Dagster as in the larger projects, once the models need a schedule rather than a command.
Stack and links
PostgreSQL, Debezium, Redpanda, ClickHouse, dbt, Docker Compose, Make, Bash and GitHub Actions. The repository, linked below, holds the code and a README that walks through how data moves, what each check catches, the design choices, and how to adapt the template to your own tables. It is MIT licensed, and a GitHub Actions workflow runs the same make demo on a fresh runner for pull requests. The two larger projects it comes from are linked below as well.