Skip to main content
Mati Data, home

Data architecture and platforms

Kiva loan CDC analytics: PostgreSQL to ClickHouse, orchestrated by Dagster

An independent portfolio project on Kiva's public loan API: PostgreSQL, Debezium into Redpanda, ClickHouse with dbt models and tests, Dagster orchestration, and Prometheus and Grafana with a reconciliation monitor.

  • 15 dbt tests and 15 unit tests
  • 4 Grafana alert rules, routed to email
  • 3 CI stages, ending in an end-to-end CDC test

Context

I built this independent portfolio project on Kiva's public loan API, which needs no key or account. It is not affiliated with or endorsed by Kiva. A Python job pulls funded loans into PostgreSQL, Debezium streams every change out of the write-ahead log into Redpanda, ClickHouse consumes the topic through its Kafka engine, and dbt builds and tests staging, intermediate and mart models, including a machine-learning feature table. Dagster runs the pipeline, Prometheus and Grafana watch the databases and the stream, GitHub Actions tests the whole path, and docker compose up -d starts all of it.

I designed it around the question every CDC pipeline is eventually asked: how do you know the connector did not quietly drop events or fall behind, when every service still reports healthy?

Constraints

  • One command. docker compose up -d brings up every service. A one-shot container renders the Debezium connector config from .env and registers it, skipping the step if the connector already exists, and Dagster runs the pipeline once at startup before its schedule takes over.
  • Laptop-sized. Every service has an explicit memory limit, so exhaustion shows up as a container killed for memory instead of a slow host, and the whole stack needs about 4 GB of free memory.
  • A free public API. The client sends a browser User-Agent because the API's firewall rejects the default Python client, retries with exponential backoff, and polls every fifteen minutes: often enough to stay current, rarely enough not to hammer a free service.
  • Everything as code: the connector config, the ClickHouse tables, the dbt models and tests, and Grafana's datasources, dashboards, alert rules, contact point and notification policy.

Architecture

  1. Kiva public APIFunded loans, no key
  2. Python ingestionRetry with backoff, upsert on the loan id
  3. PostgreSQL 16Logical replication for change capture
  4. Debezium 2.5Reads the write-ahead log through pgoutput
  5. RedpandaKafka-compatible event log
  6. ClickHouseKafka engine into ReplacingMergeTree
  7. dbtStaging, intermediate and marts
  8. GrafanaLoan analytics over the marts
Orchestration, across every stage
Dagster: ingestion, then dbt run and dbt test; GitHub Actions CI
Observability, across every stage
Prometheus and Grafana; cdc-monitor: row drift, replication lag, connector state; alert email through Mailpit
The Kiva loan CDC pipeline: data moves top to bottom, and orchestration and observability run across every stage.

Ingestion is a small Python client that fetches a few hundred funded loans per run and upserts them into PostgreSQL 16 on the loan id. A re-run updates the status, the funded amount and an updated_at timestamp on existing rows, so CDC has real updates to capture, not only inserts.

Debezium 2.5 reads the WAL through pgoutput and publishes flattened row images to Redpanda as plain JSON, with the operation and a deleted flag on every event. ClickHouse reads the topic through a Kafka engine table, and a materialized view lands each event in a ReplacingMergeTree table ordered by loan id and partitioned by month, carrying the source updated_at along as a freshness signal.

dbt builds a staging view that reads the raw table with FINAL and drops deleted rows, an intermediate view that adds a funding tier and a funding percentage, a mart by country and sector for dashboards, and a per-loan feature table with log transforms, one-hot flags, date features and a fully-funded label. The borrower name is tagged as personal data in the dbt metadata. Dagster runs two assets in order, ingestion and then dbt run followed by dbt test, once at startup and every fifteen minutes after; the CDC path itself does not wait for that schedule.

Prometheus scrapes Redpanda and ClickHouse natively, PostgreSQL through postgres-exporter, and a custom cdc-monitor exporter that polls both databases and the Kafka Connect API. Grafana is provisioned from the repository with two dashboards, one for pipeline health and one for loan analytics over the marts, plus the alert rules. dbt docs and the Redpanda and Debezium consoles run in the same stack for inspection.

Decisions and tradeoffs

Three decisions I would defend, each aimed at a failure that infrastructure metrics do not show:

  • Reconcile instead of inferring. Healthy Redpanda and ClickHouse metrics do not prove that every row arrived, so cdc-monitor compares the PostgreSQL row count with the deduplicated ClickHouse count, and measures freshness as the age of the newest replicated updated_at. A stalled or lossy connector shows up as drift or lag even while every service reports healthy.
  • No data is not zero lag. Until a row has replicated, the monitor reports no lag at all rather than a reassuring zero, so an empty warehouse cannot pass for a fresh one.
  • Connector health comes from the Kafka Connect REST API, not from throughput. The connector counts as healthy only when it and every one of its tasks report running, and a connector with no tasks counts as down, because from the outside an idle stream and a dead one look the same.

Two more are worth naming. The raw table is a ReplacingMergeTree read through FINAL in exactly one staging model, so dashboards and features see the latest state of each loan rather than every CDC version, at the cost of merge work at read time that the scaling notes plan around. And monitoring runs inside Compose instead of a hosted service, so nothing needs an external account, while Redpanda stands in for Kafka to keep the broker light enough for a laptop.

Outcome

  • 15 dbt tests, unique, not-null and accepted-value checks plus a custom singular test, and 15 unit tests for the ingestion client and the monitor's drift, lag and connector-health logic, all mocked so they run without live services.
  • 3 CI stages in GitHub Actions. Lint with unit tests, and docker compose config with dbt parse, run first. Only when both pass does an end-to-end job start PostgreSQL, Redpanda, Debezium and ClickHouse, register the connector, ingest live data from the API, wait for the rows to reach ClickHouse, run dbt run and dbt test against the warehouse, and check that the monitor reports real metrics.
  • 4 Grafana alert rules, provisioned as code, for row drift, replication lag, a stopped connector and stalled ingestion, routed through a contact point and a notification policy to email, which Mailpit catches locally. I tested the path with a real fault: stopping the Debezium container flipped the connector metric, the rule went from pending to firing, the alert email arrived, and a resolved email followed after the restart.
  • A Grafana crash traced to its cause instead of papered over. Once alerting was added, Docker began killing Grafana for running out of memory. The causes were bundled app plugins installed on every boot and a memory limit sized for dashboards alone. My first fix, a global switch, also blocked the one plugin the stack needs, which I caught in the same investigation. A one-shot installer container now adds only the ClickHouse datasource plugin, the default installs stay off and the limit was resized, and soak tests confirmed the fix.
  • Stable datasource ids. A dashboard had hard-coded a datasource id from an earlier local Grafana session, which would have blanked every panel on a fresh clone; both datasources now have fixed ids that the dashboards and alert rules share.

What I would change

The design notes carry a scaling plan by volume tier. In short: incremental dbt models and a TTL on the raw table once volume grows past demo scale. A schema registry with Avro or Protobuf instead of plain JSON, to stop upstream schema drift before it reaches ClickHouse. Debezium reading from a standby replica rather than the primary, with PgBouncer in front of PostgreSQL. At higher volumes, Kafka with partitioned topics, more connector tasks and a replicated ClickHouse cluster. Feast in front of the feature table for online serving. And with more sources and teams, Airflow with dedicated executors, a dbt project split by domain, and lineage through OpenLineage, because per-table row-count reconciliation stops being enough.

PostgreSQL, Debezium, Redpanda, ClickHouse, dbt, Dagster, Prometheus, Grafana, Docker Compose, GitHub Actions and Python. The repository, linked below, holds the code, a design report with the ClickHouse table-design rationale and the scaling plan, and notes on the data model and observability. The loan data comes from Kiva's public API; the project is not affiliated with or endorsed by Kiva.