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
Simplified architecture. Names are generic.
- 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.
- 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.
- 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.
- 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.
- 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
- Micro-batches or single events. Small batches keep latency low, bigger batches give better throughput and fewer files. You tune the batch size or interval against how fresh the data needs to be.
- Parquet or Iceberg. Iceberg handles updates properly but adds overhead, so use it only where rows actually change and keep plain Parquet for append-only data.
- Ordering. Kafka only guarantees order within a partition. Key events by the row's primary key so changes to the same row are applied in order.
- Duplicates and replays. Delivery is at-least-once, so the same event can arrive twice. Commit Kafka offsets only after the write succeeds, and make writes safe to repeat.
- Deletes. Decide whether a delete removes the row from the lake or marks it as deleted. Keeping a flag preserves history, and removing it saves space.
What to watch out for
- Lag. Track consumer lag and the delay from the source change to the lake, and alert when it grows.
- Schema changes. Columns will be added sooner or later. Plan for schema evolution so a new column doesn't break the consumers.
- Small files. Frequent micro-batches create many small files, so schedule compaction.
- The starting point. CDC gives you changes, not history. Take an initial snapshot of each table before you start streaming.
- Retention. How long Kafka keeps data decides how far back you can replay.
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.