Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Ingestion Patterns and Change Data Capture

Log-based CDC


1. The ingestion pattern menu

PatternHowCaptures deletes?Captures every change?Source loadComplexity
Full snapshotSELECT * every run✅ (by diff)❌ intermediate statesHighLow
Incremental (high-watermark)WHERE updated_at > :last_max❌❌MediumLow
Log-based CDCRead WAL/binlog (Debezium, DMS, Fivetran HVR, Datastream)✅✅ in commit orderVery lowMedium–high
Trigger-based CDCDB triggers write to an audit table✅✅High (write amplification)Medium
OutboxApp writes event to outbox table in same txn → CDC✅ (as events)✅ business eventsLowMedium
Event streamingApp publishes events directly to Kafkan/a✅n/aDual-write risk
File dropPartner drops CSV/JSON/Parquet in a bucket → AutoLoaderdependsdependsn/aLow
API pullPaginated REST/GraphQL with cursorsrarely❌Rate limitsMedium

Choosing

flowchart TD
    Q1{Source is an OLTP database<br/>you can read the log of?}
    Q1 -->|yes| Q2{Need deletes or every<br/>intermediate change?}
    Q2 -->|yes| CDC[Log-based CDC]
    Q2 -->|no, and table is small| SNAP[Full snapshot + diff]
    Q2 -->|no, large table, reliable updated_at| INC[Incremental watermark]
    Q1 -->|"no (SaaS)"| Q3{Has webhooks / change feed?}
    Q3 -->|yes| WH[Webhooks / change API → queue]
    Q3 -->|no| API[API pull with cursor + periodic full reconcile]
    Q1 -->|"files from partners"| FILES[Object storage + AutoLoader / Snowpipe]
    Q1 -->|"our own app"| OUT[Outbox pattern or events to Kafka]

2. Incremental high-watermark extraction: and its traps

-- run parameter: :last_watermark (from state table), :run_upper = now() - interval 5 minutes
SELECT * FROM orders
WHERE updated_at >  :last_watermark
  AND updated_at <= :run_upper;
-- after successful load: last_watermark = :run_upper

Traps:

  1. Deletes are invisible. Hard-deleted rows never show up. Need soft deletes or a periodic full key reconcile.
  2. updated_at not maintained by every code path (bulk updates, manual fixes) → silent misses.
  3. Long-running transactions: a row committed after your extraction but with an updated_at before your watermark is skipped forever. Mitigation: upper bound lagging now() by a safety margin, plus overlap windows with idempotent MERGE downstream.
  4. Clock/timezone issues between app servers and DB.
  5. Intermediate states lost: if status went PAID → SHIPPED between runs, you never see PAID.

3. Log-based CDC in depth

How it works

  1. The DB writes every change to its transaction log (Postgres WAL, MySQL binlog, Oracle redo, SQL Server CDC tables).
  2. A CDC connector (Debezium on Kafka Connect) registers as a replication client (Postgres: replication slot + pgoutput plugin) and streams changes.
  3. Each change becomes an event with before, after, op (c/u/d/r), source position (LSN), transaction id, timestamp.
  4. Events go to a Kafka topic per table, keyed by primary key → per-key ordering.

Applying CDC to the lakehouse

-- silver current-state table (SCD1) from a micro-batch of CDC events
MERGE INTO silver.customers t
USING (
  SELECT * FROM cdc_batch
  QUALIFY ROW_NUMBER() OVER (PARTITION BY id ORDER BY lsn DESC) = 1   -- latest change per key in this batch
) s
ON t.id = s.id
WHEN MATCHED AND s.op = 'd'                 THEN DELETE
WHEN MATCHED AND s.lsn > t._lsn             THEN UPDATE SET *        -- ignore stale/out-of-order
WHEN NOT MATCHED AND s.op != 'd'            THEN INSERT *;

Databricks DLT / Lakeflow wraps this as APPLY CHANGES INTO / AUTO CDC with SEQUENCE BY lsn and STORED AS SCD TYPE 2.

The hard parts (interview gold)

ProblemSolution
Initial load of a 2 TB tableConsistent snapshot at LSN X, then stream from X (Debezium does this; incremental snapshots avoid long locks). Overlap is safe because MERGE is idempotent on (pk, lsn).
Ordering across partitionsOnly per key is guaranteed. Multi-table transactional consistency needs txId grouping or accepting eventual consistency.
Out-of-order / duplicate deliveryOrder by LSN not timestamp; WHEN MATCHED AND s.lsn > t._lsn.
Schema changes in the sourceSchema registry with compatibility rules; Delta mergeSchema for additive changes; contract for breaking ones.
Replication slot lagIf the connector is down, Postgres retains WAL → source disk fills up → production outage. Alert on slot lag; set max_slot_wal_keep_size.
Deletesop='d' + tombstone; decide hard vs soft delete in silver; GDPR deletion must propagate to gold and backups.
TOAST / unchanged large columns (Postgres)Debezium may emit placeholders for unchanged large values → use REPLICA IDENTITY FULL or coalesce with the existing value.
High-churn tablesCompaction in Kafka; batch MERGE intervals of 1–5 min to avoid tiny commits.

4. File ingestion (AutoLoader pattern)

(spark.readStream.format("cloudFiles")
   .option("cloudFiles.format", "json")
   .option("cloudFiles.schemaLocation", "/schemas/partner_orders")
   .option("cloudFiles.schemaEvolutionMode", "addNewColumns")   # or rescue
   .option("cloudFiles.useNotifications", "true")                # queue-based discovery at scale
   .load("s3://landing/partner/orders/")
   .select("*", "_metadata.file_path", "_metadata.file_modification_time")
   .writeStream
   .option("checkpointLocation", "/chk/bronze_partner_orders")
   .trigger(availableNow=True)
   .toTable("bronze.partner_orders"))

Key points:


5. API ingestion


6. Streaming ingestion from apps (events)


7. Interview questions

Why is log-based CDC preferred over querying updated_at?

Captures deletes and every intermediate change in commit order, doesn’t depend on application discipline around updated_at, and puts negligible load on the source (reads the log instead of scanning tables). Costs: more infrastructure (Kafka Connect, registry), DB permissions/config (replication slots, binlog row format), and operational risk like WAL retention.

Debezium was down for 6 hours. What happens and how do you recover?

Postgres retained WAL for the slot (watch disk!). On restart Debezium resumes from its last committed LSN and replays 6 hours of changes; downstream MERGE by (pk, lsn) is idempotent, so consumers just catch up. Lag alerts should have fired. If the slot was dropped or WAL removed, you need a new snapshot (incremental snapshot) and reconcile.

How do you build SCD2 history from a CDC stream?

For each key, order changes by LSN; each change closes the current row (valid_to = change_ts, is_current = false) and inserts a new row (valid_from = change_ts, valid_to = '9999-12-31', is_current = true). Within a micro-batch you must handle multiple changes per key in order (window over LSN with LEAD to compute valid_to). Only track changes in columns that matter (hash of tracked columns) to avoid noise. Or use DLT APPLY CHANGES ... STORED AS SCD TYPE 2.

What's the dual-write problem and how does the outbox pattern solve it?

An app that writes to its DB and publishes to Kafka can fail between the two, leaving them inconsistent. With the outbox pattern, the app writes the business change and an event row to an outbox table in one local transaction; CDC reads the outbox and publishes it. The event is published if and only if the transaction committed, in commit order.