Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Design a CDC pipeline to the lakehouse

Pipelines & scenarios · System Design Questions

Design a CDC pipeline to the lakehouse

Hardpipelines-33
scenariocdcdebeziummysqllakehouse

Question

Design a change data capture pipeline from MySQL to a lakehouse with history.

Solution

The core idea: capture every change from MySQL's binlog, keep them all in an append-only bronze table, then derive a current-state table and a history table from that log.

MySQL binlog -> Debezium -> Kafka (topic per table)
   -> bronze (append-only change log)
   -> silver (current state, MERGE)   -> gold / SCD2 history

Capture

A connector such as Debezium reads the binlog (on GCP, Datastream does a similar job) and publishes one event per row change, with the operation (c, u, d), the before and after images, a timestamp, and the binlog file and position. Use one topic per table, keyed by the primary key, so changes to one row stay in order.

Initial snapshot, then streaming

A new table starts with a consistent snapshot of the existing rows, and then continues from the binlog position where the snapshot was taken. The connector handles that switch, but you must keep binlog retention longer than any outage, or you cannot resume and need a new snapshot.

Bronze: the change log

Land every event untouched into an append-only table (Delta or Iceberg), with the source position and ingest time. Never update it. It is your replay source.

Silver: current state

Apply changes with MERGE, but in the right order. Deduplicate the micro-batch first, keeping the latest change per key by source position (not by arrival time, which can be out of order).

MERGE INTO silver.customers t
USING (SELECT * FROM batch QUALIFY ROW_NUMBER() OVER
       (PARTITION BY id ORDER BY binlog_pos DESC) = 1) s
ON t.id = s.id
WHEN MATCHED AND s.op = 'd' THEN DELETE
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED AND s.op <> 'd' THEN INSERT *;

Decide whether deletes remove the row or set a flag. Many teams keep a soft-delete flag for audit.

History

For an SCD Type 2 table, close the old version and open a new one each time tracked columns change. The bronze log makes this repeatable, and you can rebuild history from it if the logic changes.

Problems to cover

  • Schema evolution: a schema registry for the events, and additive columns flowing into bronze, with review before they reach silver.
  • Tombstones: Kafka emits a null-valued message after a delete for compaction. Handle it, and do not treat it as a bad record.
  • Small files: streaming into a lakehouse table creates many files, so compact on a schedule.
  • Monitoring: connector lag, binlog position gap, and a daily row count reconciliation between MySQL and silver.
PreviousNext