Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Kafka for CDC with Debezium

Kafka · Operations & Scenarios

Kafka for CDC with Debezium

Mediumkafka-59
cdcdebeziumdatabase-replicationkafka-connect

Question

How does Debezium capture database changes into Kafka?

Solution

Debezium captures database mutations by reading low-level database transaction logs (such as the PostgreSQL write-ahead log via logical replication or the MySQL binary log) and streaming every row-level change to dedicated Kafka topics. It executes an initial consistent snapshot of existing tables before switching to real-time log tailing, emitting structured change events containing before and after row states, operation types, and transaction positions. By publishing each source table to its own topic and partitioning records by table primary key, Debezium preserves strict chronological ordering for every database record.

Tailing the transaction log

Unlike query-based polling that stresses databases and misses intermediate row updates, Debezium connects as a replication follower:

  • In PostgreSQL, it uses logical decoding plugins to read the Write-Ahead Log (WAL).
  • In MySQL, it reads the row-based binary log (binlog).

Because it reads committed transaction logs directly, Debezium incurs minimal CPU overhead on the primary database, captures intermediate updates that occurred between polls, and reliably catches row deletions.

Structure of a Debezium change event

Each emitted message contains comprehensive envelope metadata:

  • op: The mutation operation (c for create, u for update, d for delete, r for initial read snapshot).
  • before: The row column values prior to the change (null on inserts).
  • after: The row column values after the change (null on deletes).
  • source: Transaction metadata including database name, table name, log position, and commit timestamp.

When a row is deleted, Debezium produces a delete record with the previous state, immediately followed by a tombstone message (a record with the same key but a null value), enabling log-compacted Kafka topics to purge the key.

Snapshots and primary key partitioning

When starting on an existing database, Debezium performs an initial snapshot of all historical rows without acquiring write locks, transitioning to streaming at the exact transaction offset where the snapshot ended. Messages are partitioned using the table primary key, ensuring all sequential changes for a specific database row land on the exact same Kafka partition in chronological order.

PreviousNext