Command Palette

Search for a command to run...

Hectal
PHASE 9Intermediate ~9 min· topic 3 of 4

Topic 9.3

Change Data Capture with Debezium

In one line

CDC reads a database's transaction log (PostgreSQL WAL via logical replication, MySQL binlog) and emits every insert, update and delete as an event, in commit order per row, starting with an initial snapshot. It's the reliable way to stream database changes to Kafka without changing application code.

0/4 · 0%

Think of it like this

A court stenographer. Instead of asking every lawyer to report what they said, the stenographer records everything said in order from the official transcript. CDC records every committed change from the database's own log.

Key ideas

  1. 01

    How Debezium PostgreSQL works: a logical replication slot (with the pgoutput plugin) streams committed changes; the connector converts them to change events with before, after, op (c, u, d, r for snapshot reads), source metadata (LSN, transaction ID, timestamp).

  2. 02

    Initial snapshot: on first start, Debezium reads existing rows (op r) and then streams changes from the log position captured at snapshot time, so consumers get a full copy then increments. Incremental snapshots (signalling table) re-snapshot tables without stopping streaming.

  3. 03

    Keys and ordering: events are keyed by the table's primary key and go to one topic per table (server.schema.table), so changes to one row stay ordered. There's no ordering across tables or rows unless you route them together.

  4. 04

    Deletes: a delete event (op d) followed by a tombstone (null value) for log-compacted topics, so compaction eventually removes the key.

  5. 05

    Schema changes: column additions flow through (with Schema Registry, as new schema versions); breaking changes need coordination. Transaction boundaries can be emitted as separate events if consumers need them.

  6. 06

    Operational risk: a replication slot that no one consumes makes PostgreSQL retain WAL, which can fill the disk. Monitor slot lag and set max_slot_wal_keep_size.

Code & diagrams

debezium-postgres.jsonjson
{
  "name": "commerce-cdc",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "pg.internal",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "${file:/secrets/pg.properties:password}",
    "database.dbname": "commerce",
    "topic.prefix": "cdc.commerce",
    "plugin.name": "pgoutput",
    "slot.name": "debezium_commerce",
    "publication.autocreate.mode": "filtered",
    "table.include.list": "public.orders,public.customers",
    "snapshot.mode": "initial",
    "tombstones.on.delete": "true",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter.schema.registry.url": "http://registry:8081",
    "value.converter.schema.registry.url": "http://registry:8081"
  }
}
change-event.jsonjson
{
  "before": { "id": 9, "status": "CREATED", "total": 1298.00 },
  "after":  { "id": 9, "status": "PAID",    "total": 1298.00 },
  "op": "u",
  "ts_ms": 1727520005331,
  "source": { "connector": "postgresql", "db": "commerce", "table": "orders",
              "lsn": 24023128, "txId": 5521, "snapshot": "false" }
}
slot-lag.sqlsql
SELECT slot_name, active,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS lag
FROM pg_replication_slots;
--    slot_name       | active |  lag
-- debezium_commerce  | t      | 12 MB

Interview problem

The problem

Database CDC to many consumers

Design PostgreSQL → CDC → Kafka → multiple consumers (search, cache invalidation, analytics). Cover ordering, deletes, schema changes, the initial snapshot, incremental updates, exactly-once limitations and replay.

The interviewer follows up

01

Why not have the application publish events itself instead of CDC?

When it breaks

The CDC connector is stopped for days and its replication slot stays

What you see

PostgreSQL keeps all WAL since the slot's position; the database disk fills and the primary goes down.

Fix & prevent

Alert on slot lag, set max_slot_wal_keep_size, and drop slots of decommissioned connectors.

Explain it without notes

01

What does a Debezium change event contain, and how are deletes represented?

Practice

01

Run Postgres + Debezium in Docker, insert, update and delete a row, and consume the change topic.

Trade-offs

  • ↔

    CDC captures every change reliably but exposes table structure, coupling consumers to the database schema.

Done when you can

  • I can design a Debezium CDC pipeline with snapshots, deletes, schema changes and replay.