skip to content

How does a Debezium CDC pipeline move database changes through Kafka into an analytics lakehouse, and what makes it preferable to periodic batch extracts?

level: seniorimportance: must knowfreq 60%

answer

  1. log-based: binlog / pgoutput / LogMiner
  2. event: before/after + op c/u/d/r + source LSN
  3. snapshot then stream handoff
  4. tombstone for deletes; PK upsert/merge
  5. vs batch: fresh, low load, captures deletes

basics

~20 s

Debezium is a Kafka Connect source connector that reads a database's transaction log and emits each insert/update/delete as a Kafka event. Sink connectors then land those change events in the lake, giving near-real-time, low-load, complete change history instead of nightly full-table dumps.

solid answer

~50 s

Debezium runs as a Kafka Connect **source connector** that performs **log-based CDC**: it tails the database's write-ahead/transaction log (Postgres logical decoding/`pgoutput`, MySQL binlog, Oracle LogMiner, SQL Server CDC) and emits a structured change event per row change to a per-table Kafka topic. Each event carries `before`/`after` images, `op` (c/u/d/r), source metadata (LSN, txid, timestamp), and is typically Avro+Schema Registry encoded. It starts with a consistent **initial snapshot**, then streams incrementally. Downstream, a sink (S3/Iceberg/BigQuery/Snowflake) lands the stream; analytics either keeps the append-only change log or materializes the current state via **upsert/merge** keyed on the primary key, handling deletes via tombstones. Versus nightly batch extracts, CDC gives near-real-time freshness, minimal source load (reads the log, not the tables), captures *every* change including intermediate states and deletes, and avoids brittle high-watermark queries. Key concerns: exactly-once/dedup, schema evolution, and snapshot-to-stream handoff.

go deeper

for a junior

Know Debezium turns DB inserts/updates/deletes into Kafka events for near-real-time pipelines.

for a middle

Explain log-based capture, the change-event structure, and snapshot-then-stream.

for a senior

Detail materialized upserts, tombstones/deletes, schema evolution, and why CDC beats high-watermark batch.

for a principal

Architect exactly-once/idempotent merge, partitioning by PK for ordering, and operational concerns (slot bloat, snapshot strategy).

## What CDC is **Change Data Capture (CDC)** means capturing every row-level change (insert/update/delete) in a source database as a stream of events, rather than periodically re-reading whole tables. **Log-based CDC** reads the database's own **transaction log** — the durable, ordered record the DB already writes for crash recovery and replication. ## Debezium's role **Debezium** is a set of **Kafka Connect source connectors**. For each supported DB it taps the log: - **Postgres**: logical decoding via the `pgoutput` plugin and a replication slot. - **MySQL/MariaDB**: the **binlog**. - **Oracle**: LogMiner (or XStream). - **SQL Server**: the built-in CDC tables. It produces one Kafka topic per table by default, with a structured **change event** per row change containing: - `before` and `after` row images, - `op`: `c` (create/insert), `u` (update), `d` (delete), `r` (read, from snapshot), - `source` block: LSN/binlog offset, transaction id, commit timestamp, table name. Events are usually serialized as **Avro with Schema Registry** so schema changes propagate safely. ## Snapshot then stream When first started, Debezium takes a **consistent initial snapshot** of existing rows (emitted as `op=r`) so the lake has the full baseline, then switches to **streaming** from the log position captured at snapshot time — a careful handoff so no change is missed or duplicated at the boundary. **Incremental snapshots** (signal-based, watermark windows) let you re-snapshot tables without stopping streaming. ## Landing in the lakehouse A sink connector writes the change stream to the lake. Two common shapes: 1. **Append-only change log**: store every event; analysts reconstruct state or do temporal queries. Great for audit and time-travel. 2. **Materialized current state**: an **upsert/merge** keyed on the primary key applies inserts/updates and removes rows on delete. Debezium emits a **tombstone** (null value) after a delete so log-compacted topics and merge logic can drop the key. Iceberg/Snowflake/BigQuery sinks support MERGE-style upserts; the Debezium-aware `ExtractNewRecordState` SMT flattens the envelope to just the `after` image for simpler sinks. ## Why CDC beats batch extracts - **Freshness**: seconds/minutes vs. hours; powers near-real-time dashboards and ELT. - **Low source load**: reads the log, not the live tables — no heavy full scans, no lock contention. - **Completeness**: captures *every* change, including **intermediate updates and deletes** that a high-watermark `WHERE updated_at > X` query silently misses (and that query can't see hard deletes at all). - **No watermark fragility**: batch extracts depend on a reliable monotonic column and miss late/out-of-order updates. ## Edge cases to name - **Duplicates/exactly-once**: Connect is at-least-once; the lake merge must be idempotent (keyed upsert) or use dedup on LSN/offset. - **Schema evolution**: column adds/drops flow through Schema Registry; downstream tables must evolve compatibly. - **Deletes**: handle tombstones; otherwise deleted rows linger in materialized tables. - **Ordering**: per-key ordering holds within a partition; key topics by PK to preserve update order. Debezium-through-Kafka is the de facto pattern for feeding fresh, complete operational data into the analytics lakehouse.

  • Why does a high-watermark batch extract (WHERE updated_at > last_run) miss data that CDC captures?
    It can't see hard deletes (the row is gone, no updated_at to catch), misses intermediate update states between runs, and breaks if updated_at isn't reliably monotonic or if rows are updated without bumping it. CDC reads the log so it sees every committed change in order.
  • How are deletes represented end-to-end so a materialized lake table actually drops the row?
    Debezium emits a change event with op=d (and the before image), followed by a tombstone record (null value) for the key. The sink's upsert/merge logic interprets the delete to remove the keyed row; on compacted topics the tombstone lets Kafka physically drop the key.

saying these in an interview costs you the question

  • Describing Debezium as query/polling-based — it is log-based CDC reading the transaction log.
  • Forgetting deletes/tombstones, leaving deleted rows in the materialized table.
  • Assuming the pipeline is exactly-once — Connect is at-least-once; the merge must be idempotent.
  • Claiming CDC adds heavy load to source tables — it reads the log, not the tables.

context