Skip to content
Reliable Data Engineering
Practice problem medium cdcdebeziumdeltascd2merge
Practise with timer, notes and rubric

Design a CDC Pipeline from 200 OLTP Tables into a Lakehouse

CDC pipeline

Problem

The company runs its business on several Postgres and MySQL databases (~200 tables, the largest 3 TB). Analytics currently runs nightly SELECT * dumps that take 6 hours and load the production DBs. Design a replacement that delivers fresh, historised data to the lakehouse.


Clarifying questions

QuestionAssumed answer
Freshness needed?≤ 15 min for most tables, ≤ 1 min for orders/payments
Change volume?~50 M row changes/day overall, bursts during batch jobs in the source (10× for 1 h)
History needed?Yes: point-in-time (“what was the customer’s tier on date X?”)
Deletes?Hard deletes happen in sources; must be reflected; GDPR deletes must propagate
Schema changes?Weekly app deploys add columns; occasional renames
Can we change the source DBs?We can enable logical replication / binlog; no app code changes

1. Requirements

Low source impact; capture inserts/updates/deletes; current-state tables (SCD1) and history (SCD2); handle schema evolution; replay/backfill; monitoring of lag per table; onboarding a new table is config, not code.

2. Estimates

50M changes/day ≈ 600/s avg, 6k/s burst × ~1 KB = 6 MB/s → small for Kafka
Initial snapshot: ~10 TB total across DBs → incremental snapshots, throttled, over days
Silver MERGE cadence: 1–5 min micro-batches; history tables grow ~50M rows/day

3. Architecture

flowchart LR
    subgraph Sources
        PG[(Postgres<br/>WAL + replication slot)]
        MY[(MySQL<br/>binlog ROW format)]
    end
    PG --> DBZ[Debezium connectors<br/>Kafka Connect cluster]
    MY --> DBZ
    DBZ --> SR[Schema Registry]
    DBZ --> K[[Kafka: one topic per table<br/>key = primary key]]
    K --> ING[Generic ingest job<br/>config-driven, one per DB]
    ING --> BR[(bronze.table_cdc<br/>append-only change log)]
    BR --> CUR[MERGE → silver.table<br/>current state]
    BR --> HIS[SCD2 → silver.table_history]
    CUR --> GOLD[(gold models via dbt)]
    HIS --> GOLD
    CFG[(Table config:<br/>keys, SCD type, freshness tier,<br/>PII columns)] --> ING
    CFG --> CUR

4. Data model

-- bronze change log: one row per change event (never updated)
bronze.orders_cdc (
  op STRING,                -- c, u, d, r (snapshot)
  pk STRING,                -- serialized primary key
  after STRUCT<...>, before STRUCT<...>,
  lsn BIGINT, tx_id BIGINT, source_ts TIMESTAMP,
  _kafka_offset BIGINT, _ingested_at TIMESTAMP
) PARTITIONED/CLUSTERED BY (ingest_date)

-- silver current state (SCD1): mirrors source + metadata
silver.orders (... source columns ..., _lsn BIGINT, _updated_at TIMESTAMP, _is_deleted BOOLEAN)

-- silver history (SCD2)
silver.orders_history (... columns ..., valid_from TIMESTAMP, valid_to TIMESTAMP, is_current BOOLEAN, _lsn BIGINT)

5. Deep dives

5.1 Config-driven generic pipeline

Writing 200 pipelines by hand doesn’t scale. One parameterised job per database reads all its topics, with per-table config:

- table: public.orders
  primary_key: [order_id]
  scd: [1, 2]                 # current + history
  freshness_tier: 1min
  sequence_by: lsn
  track_history_columns: [status, amount, shipping_address_id]   # SCD2 only on meaningful columns
  pii: [customer_email]
  soft_delete: true

Onboarding a table = PR to config + connector table include list. CI validates the config (keys exist, etc.).

5.2 Applying changes correctly

MERGE INTO silver.orders t
USING (
  SELECT after.*, op, lsn FROM batch_changes
  QUALIFY ROW_NUMBER() OVER (PARTITION BY pk ORDER BY lsn DESC) = 1
) s
ON t.order_id = s.order_id
WHEN MATCHED AND s.lsn <= t._lsn THEN UPDATE SET t._lsn = t._lsn   -- stale event, no-op (or filter out beforehand)
WHEN MATCHED AND s.op = 'd'      THEN UPDATE SET _is_deleted = true, _lsn = s.lsn, _updated_at = current_timestamp()
WHEN MATCHED                     THEN UPDATE SET *
WHEN NOT MATCHED AND s.op != 'd' THEN INSERT *;

5.3 SCD2 from a change stream

-- 1) Within the micro-batch, give every change a validity interval (per key, in LSN order)
WITH ordered AS (
  SELECT pk, after.*, op, lsn,
         source_ts                                                   AS valid_from,
         LEAD(source_ts) OVER (PARTITION BY pk ORDER BY lsn)         AS next_ts
  FROM batch_changes
)
SELECT *,
       COALESCE(next_ts, TIMESTAMP '9999-12-31')                     AS valid_to,
       next_ts IS NULL                                               AS is_current
FROM ordered
WHERE op != 'd';     -- a delete creates no new version; it only closes the previous one (step 2)

-- 2) Close the currently open row in silver.orders_history for each pk in the batch:
--    valid_to = MIN(source_ts) of that pk's changes in this batch, is_current = false
-- 3) Insert the rows from step 1.
-- Skip changes where the hash of tracked columns equals the previous version's hash (no-op updates).

In Databricks this whole section is AUTO CDC ... SEQUENCE BY lsn STORED AS SCD TYPE 2 TRACK HISTORY ON (status, amount). Know both the declarative and the hand-written version.

5.4 Initial load and backfill

5.5 Schema evolution

Source changeHandling
Add columnRegistry accepts (backward compatible); ingest uses mergeSchema; silver gets new nullable column automatically
Drop columnKeep column in silver (nulls going forward); alert consumers
RenameAppears as drop + add → detect via config/schema diff alerts; map explicitly
Type changeWidening OK; narrowing/incompatible → quarantine + alert, manual migration

5.6 Monitoring

Per table: replication slot lag (bytes) at source, connector status, Kafka consumer lag, now() − max(source_ts) in silver (freshness by tier), row count reconciliation vs source (daily COUNT(*) on a replica), MERGE duration, quarantine counts.

6. Trade-offs

DecisionChoiceAlternative
Capture methodLog-based (Debezium)Managed (Fivetran/DMS/Datastream): less ops, more cost and less control; query-based misses deletes
Kafka in the middleYesDirect connector → Delta (e.g. Debezium Server, managed CDC): fewer parts, but loses replay and fan-out
One job per DB vs per tablePer DB (multiplexed)Per table: isolation but 200 clusters; middle ground: per freshness tier
HistorySCD2 in silver for selected tables + Delta CDF for othersSnapshot tables daily: simpler but coarse and storage-heavy

7. Failure modes

FailureHandling
Connector down → WAL accumulates on sourceAlert on slot lag; max_slot_wal_keep_size; runbook; worst case drop slot + re-snapshot
Poison message (unparseable)DLQ topic per connector, alert
MERGE conflicts (concurrent writers)Single writer per table; or partition-aware merges
Source failover to a new primaryDebezium resumes on new primary (slot must exist there: logical replication failover slots / re-snapshot)
Rogue backfill in source updates 100M rowsBurst absorbed by Kafka; silver lag temporarily; autoscale ingest

8. What separates a senior answer

9. Follow-up questions

How do you keep foreign-key consistency across tables in silver (orders without customers)?

CDC streams per table are independent, so momentary inconsistency is normal. Options: tolerate it in silver and handle in gold (late-arriving dimension → placeholder/“unknown” member, fixed on next run), process related tables in the same micro-batch with tx_id awareness for strict needs, or build gold on a slight delay (e.g. watermark of 5 minutes) so related changes have landed.

How would you answer "what did this order look like last Tuesday 10:00"?

Query silver.orders_history: WHERE order_id = X AND valid_from <= ts AND ts < valid_to. Or Delta time travel on silver.orders if within retention (but SCD2 is the durable, intentional history; time travel is for recovery).

Fivetran vs self-managed Debezium?

Managed: fast to start, connectors maintained, no Kafka ops; costs scale with rows (expensive at high volume), less control over latency/format. Debezium: cheaper at scale, flexible, Kafka fan-out, but you own operations (Connect cluster, upgrades, slot management). Many companies start managed and move hot/high-volume sources to self-managed.


Self-assessment rubric