CDC Pipelines with Kafka, Debezium and Snowflake

Every data platform eventually hits the same wall: the data everyone wants to see lives inside an operational database, and nobody wants analytics queries running against it. So you build a nightly batch export instead. That solves the “don’t hammer the OLTP database” problem and immediately creates a new one – every dashboard is now a day behind reality.

On a large enterprise data platform I worked on, the fix was a real-time change data capture (CDC) pipeline: every insert, update and delete in a source PostgreSQL database streamed out within seconds and landed in Snowflake, feeding live, fleet-wide analytics dashboards for business stakeholders instead of yesterday’s numbers. This post is that pipeline – Kafka, Debezium and the Snowflake Kafka Connector – plus the lessons that only show up once it’s running in production.

Why CDC instead of another nightly batch job

Batch ETL is simple and that simplicity is real value – don’t reach for CDC just because it’s trendy. But it breaks down in a specific way: the gap between “when something happened” and “when you can see it” grows exactly as large as your batch window, and that gap is where stale decisions get made. CDC closes it without asking you to rewrite the source application as event-sourced. The source system keeps writing to Postgres exactly as before; a connector reads the database’s own write-ahead log and turns every row change into an event. Nothing in the application changes.

The architecture, end to end

The pipeline has two halves, and it’s worth keeping them conceptually separate because they fail independently.

The write path – Postgres to Kafka. PostgreSQL is configured for logical replication. Debezium, running as a Kafka Connect source connector, attaches to a replication slot and turns every committed change into a structured event – before-image, after-image, operation type, source transaction metadata – published to one Kafka topic per table.

The read path – Kafka to Snowflake. The Snowflake Kafka Connector (a sink connector, also running in Kafka Connect) subscribes to those topics and writes the raw events into Snowflake landing tables as semi-structured VARIANT rows. From there, SQL does the rest – merging raw events into a clean, current-state table that looks exactly like the source table, just a few seconds behind instead of a day.

Turning on logical replication and configuring Debezium

PostgreSQL doesn’t emit a change stream by default – logical replication has to be switched on:

# postgresql.conf
wal_level = logical
max_wal_senders = 10
max_replication_slots = 10

Then the Debezium connector itself, submitted to Kafka Connect as JSON:

{
  "name": "fleet-db-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "fleet-db.internal",
    "database.port": "5432",
    "database.user": "debezium_replicator",
    "database.password": "${file:/secrets/debezium.properties:db.password}",
    "database.dbname": "fleet",
    "topic.prefix": "fleet",
    "table.include.list": "public.leases,public.vehicles",
    "plugin.name": "pgoutput",
    "slot.name": "fleet_cdc_slot",
    "publication.autocreate.mode": "filtered"
  }
}

A couple of details matter more than they look: plugin.name: pgoutput uses Postgres’s built-in logical decoding (no extra extension to install and maintain), and scoping table.include.list tightly keeps the replication slot – and the Kafka topics downstream – from filling up with tables nobody reads.

Landing change events in Snowflake

The Snowflake Kafka Connector is configured separately, pointing at the same topics:

name=snowflake-sink
connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector
topics=fleet.public.leases,fleet.public.vehicles
snowflake.url.name=https://<account>.snowflakecomputing.com
snowflake.user.name=cdc_loader
snowflake.private.key=${file:/secrets/snowflake.properties:private.key}
snowflake.database.name=RAW
snowflake.schema.name=CDC_LANDING
buffer.count.records=10000
buffer.flush.time=10

Each change event lands as a single VARIANT row with the full Debezium payload intact – including __op (c/u/d for create/update/delete) and the source transaction’s LSN, which turns out to be the single most useful field in the whole pipeline.

From raw change events to a usable current-state table

Nobody wants to query a VARIANT column full of nested JSON, so a scheduled task merges new events into a proper table:

MERGE INTO CDC_LANDING.LEASES_CURRENT AS target
USING (
    SELECT
        payload:after:id::NUMBER         AS id,
        payload:after:status::STRING     AS status,
        payload:after:updated_at::TIMESTAMP AS updated_at,
        payload:op::STRING               AS op,
        payload:source:lsn::NUMBER       AS lsn
    FROM CDC_LANDING.LEASES_RAW
    WHERE _ingested_at > DATEADD(minute, -15, CURRENT_TIMESTAMP())
    QUALIFY ROW_NUMBER() OVER (PARTITION BY payload:after:id ORDER BY payload:source:lsn::NUMBER DESC) = 1
) AS source
ON target.id = source.id
WHEN MATCHED AND source.op = 'd' THEN DELETE
WHEN MATCHED AND source.lsn > target.lsn THEN UPDATE SET
    status = source.status, updated_at = source.updated_at, lsn = source.lsn
WHEN NOT MATCHED AND source.op != 'd' THEN INSERT
    (id, status, updated_at, lsn) VALUES (source.id, source.status, source.updated_at, source.lsn);

The QUALIFY clause and the LSN comparison are doing the real work here – more on why below.

What production teaches you that the docs don’t

  • Order is not guaranteed, but LSN is. Kafka preserves order within a partition, not across retries or partition rebalances. Never trust wall-clock or ingestion order to decide which event wins – always compare the source LSN (or Debezium’s transaction offset) and let the highest one win, exactly as the MERGE above does.
  • Watch the replication slot, not just Kafka Connect. If the Debezium connector goes down, Postgres keeps the replication slot open and keeps accumulating WAL to replay later. A connector outage that looks harmless from Kafka Connect’s dashboard can quietly fill up disk on the database server. Alert on slot lag, not just consumer lag.
  • Schema changes need a process, not a hope. Debezium happily emits a new schema the moment a column is added or renamed upstream. Snowflake does not automatically evolve the target table. Agree a migration step – even a manual one – before anyone adds a column to a CDC’d table.
  • A few noisy tables can dominate your cost. One table with a high-frequency updated_at touch on every read can generate more events than everything else combined. Filter at the Debezium connector (table.include.list, column exclusions) rather than downstream in Snowflake – it’s cheaper to not capture an event than to capture and discard it.

When CDC earns its complexity (and when it doesn’t)

CDC is worth the operational overhead – a Kafka cluster, Kafka Connect, replication slot monitoring – for tables that feed live dashboards, alerting, or anything where a stakeholder would notice a day’s staleness. For reference data that changes weekly, a scheduled batch job is still the right answer, and simpler is a feature. Reach for CDC when the business cost of staleness is higher than the engineering cost of running it.

Streaming every change the moment it happens feels like overkill right up until someone asks why their dashboard doesn’t match what just happened in production. Get the plumbing right once, and that question stops coming up.

Elmo Yeldo is a data and platform engineer with 13+ years across fintech, automotive and aviation, currently building cloud-native data infrastructure (AWS CDK, Terraform, Snowflake, Airflow) for a large enterprise data platform.

Need a hand with a CDC pipeline, a Snowflake build-out, or a self-service data platform of your own? Get in touch – I help teams design and ship this kind of infrastructure.

Similar Posts

Leave a Reply

Your email address will not be published. Required fields are marked *