← Back to portfolio

Data engineering · 5 min read

Streaming database changes into a lakehouse with CDC

Reading straight from a production database is convenient, and production databases do not like it. Change Data Capture (CDC) is a well-known fix: instead of repeatedly asking the database for data, you listen to the changes it already records and stream them somewhere else. This is a general overview of one common way to build that.

The problem

Reports and other consumers often query the same transactional database that serves the application. Heavy scans compete with real traffic, queries slow down for everyone, and every new consumer adds more load. Copying the data over on a nightly schedule helps, but then every consumer works with yesterday's data.

The pattern at a glance

Source databaseCDC enabled Debeziumin a Docker container Kafkaone topic per table Python consumerson EC2, micro-batches Parquet on S3append-only tables Iceberg on S3tables with updates Downstreamconsumers

Simplified architecture. Names are generic.

  1. Capture changes at the source. Turn on CDC in the source database (SQL Server has a built-in feature for this). It records every insert, update and delete for the tables you choose. A Debezium connector, running in a Docker container as a Kafka Connect connector, reads those change records, so your applications never run extra queries against the business tables.
  2. Publish to Kafka, one topic per table. Each change becomes an event on a topic named after its table. Separate topics keep streams isolated, let different consumers read at their own pace, and make it possible to replay from an earlier point.
  3. Consume with Python and write in micro-batches. Python consumers running on an EC2 instance read the topics, group events into small batches by time or by count, and write each batch to the lake. Batching gives you fewer, larger files and far better throughput than writing event by event, at the cost of a little latency.
  4. Pick the storage format by how the table changes. Append-only data, such as event logs, goes to Parquet files on S3: simple, fast and cheap. Data that gets updated goes to Apache Iceberg tables, which support upserts (MERGE) and give you a consistent, versioned view of the table.
  5. Read from the lake, not the database. Downstream consumers query the lake, so the production database can concentrate on its own job.

What a change event looks like

Debezium wraps every change in an envelope. Here is a simplified version (not valid JSON, since it has comments):

{
  "op": "u",          // c = create, u = update, d = delete
  "ts_ms": "...",     // when the change happened
  "before": { ... },  // the row before the change (empty for inserts)
  "after":  { ... },  // the row after the change (empty for deletes)
  "source": { "table": "...", "position": "..." }
}

The op field tells the consumer what to do: append the row, merge it into the table, or apply a delete.

Design choices and trade-offs

What to watch out for

Wrapping up

CDC turns the database's own record of changes into a stream you can build on. The details vary with every stack, but the ideas carry across most designs: capture from the database's change records, stream by table, batch the writes, choose storage by change pattern, and read from the lake.

This is a general, simplified overview of a common pattern. It doesn't describe any employer's internal systems.