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.
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
- 01
How Debezium PostgreSQL works: a logical replication slot (with the
pgoutputplugin) streams committed changes; the connector converts them to change events withbefore,after,op(c,u,d,rfor snapshot reads),sourcemetadata (LSN, transaction ID, timestamp). - 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. - 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. - 04
Deletes: a delete event (op
d) followed by a tombstone (null value) for log-compacted topics, so compaction eventually removes the key. - 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.
- 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
{
"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"
}
}{
"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" }
}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 MBInterview 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
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
What does a Debezium change event contain, and how are deletes represented?
Practice
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.