Skip to main content
Mati Data

Data architecture and platforms

wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command

An independent reference implementation of a CDC analytics pipeline: a public REST API into PostgreSQL, Debezium into Redpanda, ClickHouse with dbt marts, Airflow, Prometheus and Grafana, started with one command.

  • 2,970 rows
  • 36 paginated requests
  • about 3 seconds from PostgreSQL commit to a queryable ClickHouse row
  • 58 dbt tests and 59 unit tests
  • 10 Prometheus alert rules

Context

Most of my change-data-capture work lives in private employer repositories, so I built a public reference implementation of how I think a CDC analytics pipeline should be put together and, more to the point, how it should fail. It ingests a public REST API, the World Bank Indicators API, into PostgreSQL, streams every change out of the write-ahead log with Debezium into Redpanda, lands the events in ClickHouse, and builds a staging layer, an analytics mart and a machine-learning feature table with dbt. Airflow orchestrates, Prometheus and Grafana watch, GitHub Actions tests, and make demo starts all of it.

I optimised for decisions I can defend and for failures that announce themselves. Most of what I have debugged in production was not a component that crashed but one that kept running while quietly doing the wrong thing.

Constraints

  • One command, idempotent: a second make demo exercises the change-detection and incremental paths instead of repeating the first run.
  • Laptop-sized: every image pinned, memory capped explicitly, ports bound to loopback on non-default numbers so the stack cannot collide with a database already running on the machine.
  • Scope written down before code: what is built, what is deliberately left out and why. No Schema Registry, no clustering, no history tables, no Alertmanager routing, no BI layer.
  • Every stage verifiable by a script with a meaningful exit code, because prose drifts and a script does not.

Architecture

Ingestion is a small Python client with pagination and retry with jitter, behind a contract that checks the shape of the response rather than its status code: this API returns errors inside successful responses and some transient failures as client errors. Rows load into PostgreSQL 17 through a change-detecting upsert, and the data, the audit row and the watermark commit in one transaction, so a recorded position can never claim work that rolled back.

Debezium 3.6.1 reads the WAL through pgoutput over an explicit publication and writes flattened events to Redpanda, with the operation, the LSN and a deleted flag carried alongside the row. ClickHouse reads each topic through a Kafka engine table with one raw JSON column. Materialized views do the typing and write twice: to an immutable event log, and to ReplacingMergeTree landing tables versioned by the source LSN, so the newest row follows the source's commit order rather than arrival order. dbt builds staging views that apply FINAL and the tombstone filter from exactly one macro, then two dimensions, an incremental fact and a feature table. Airflow runs the pipeline with quality gates in the critical path plus a separate liveness DAG. Prometheus scrapes ClickHouse and Redpanda natively, one exporter covers the signals that live in tables, and Grafana is provisioned from the repository.

Decisions and tradeoffs

Three decisions I would defend, each preventing a failure that produces no error at all:

  • The natural key is the primary key. Under the default replica identity a delete event carries only the primary key, so a surrogate key would make deletes arrive downstream as an integer with no way to identify the business entity. Full replica identity would fix that by multiplying WAL volume; the key choice fixes it for free.
  • The upsert is change-detecting. Without WHERE source_hash IS DISTINCT FROM EXCLUDED.source_hash, every re-ingest rewrites every row and each no-op update emits a CDC event, so the change stream describes the scheduler instead of the data.
  • The Kafka engine reads raw JSON and typing happens in the views. Typed Kafka-engine columns look cleaner and are a trap: one parse failure fails the block, offsets are never committed and the consumer retries forever with no error reaching anyone. With raw JSON the consumer cannot fail, and the pipeline degrades to nulls instead of stalling.

Two more are worth naming. CDC lag is measured against a Debezium heartbeat, not a business table, because an idle table reports the same number as a stopped connector. And Great Expectations is substituted, not skipped: dbt tests and pytest cover the relational checks, explicit assertions guard the boundary before data lands, and the design notes say where I would add Great Expectations first, at the feature boundary where the useful checks are distributional. Redpanda instead of Kafka with ZooKeeper, the Kafka engine instead of a Connect sink, and plain JSON instead of a Schema Registry each keep the stack small and inspectable, and each cost is written down.

Outcome

  • 2,970 rows ingested across 36 paginated requests; a second run writes nothing.
  • Latency of about 3 seconds from PostgreSQL commit to a queryable ClickHouse row.
  • 58 dbt tests and 59 unit tests green in CI, and 10 Prometheus alert rules, each with a runbook and each unit-tested with promtool. An audit of that claim found rules with no test at all while promtool still reported success, because it only runs the cases it is given, so the alert names in the tests are now checked against the rules defined.
  • A convention gate in CI for rules that fail silently when broken: FINAL only through the macro, every Replacing engine declaring its sort key, every incremental model declaring its unique key, pinned images, loopback-only ports, and every alert with a duration and a runbook.
  • Updates and deletes proven to propagate end to end by an executable, not asserted in prose.

What I would change

A Schema Registry with Avro, the one omission with a real correctness consequence. Replicated ClickHouse with Keeper, with the migration path and its trigger volume already written down. Alertmanager routing on the existing rules. Great Expectations at the feature boundary. dbt snapshots over the retained event log to rebuild history, which the current-state marts deliberately do not keep. Airflow runs in a demo topology here; the production topology is described in the design notes rather than half-built.

PostgreSQL, Debezium, Redpanda, ClickHouse, dbt, Apache Airflow, Prometheus, Grafana, Docker Compose, GitHub Actions, Python and Make. The repository, linked below, holds the code, the measured notes on the source API and the CDC wire format, and the CI evidence.