Skip to main content
Mati Data, home

Data architecture and platforms

CDC Starter: PostgreSQL to ClickHouse CDC in one command

A template of the CDC pattern I use, small enough to read in one sitting: PostgreSQL, Debezium into Redpanda, ClickHouse and dbt in one command, with deletes, ordering, lag and reconciliation handled.

In plain words: A small copy of my pipeline pattern that anyone can start with one command: change a few rows in a database, watch the analytics warehouse catch up, and let the checks show that nothing was lost.

  • 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 demo passes 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, make and curl. 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

  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.

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 ReplacingMergeTree hides it under FINAL. Without this, a delete either vanishes or comes back as a row with empty columns.
  • The current-state read lives in one macro. Forgetting FINAL gives 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 down removes everything, the slot included.

Outcome

  • make demo passes from a clean start, after make down has 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 FINAL and 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.

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.

Working on something like this?

Where I would start: Team workshop

Get in touch