ProjectsProject 6 of 8
Advanced project · Project 6 of 8
Change Data Capture Pipeline
Replicate an operational PostgreSQL table into a lakehouse table within minutes, including updates and deletes, so analysts query current data without touching the production database.
Requirements
- Run PostgreSQL with logical replication enabled
- Capture changes with a log-based CDC connector into Kafka
- Apply changes to a Delta table with MERGE, ordered by source log position
- Handle deletes
- Take an initial snapshot without missing concurrent changes
Technology stack
PostgreSQL, A log-based CDC connector (for example Debezium), Kafka, Spark Structured Streaming, Delta Lake, Docker Compose
Dataset
Create your own orders table and a small script that inserts, updates and deletes rows continuously.
Business context
Analysts need current operational data, but querying the production database directly risks slowing the application. CDC keeps an analytical copy up to date from the database’s own change log, without extra query load.
Architecture
- PostgreSQL writes changes to its write-ahead log.
- The CDC connector reads the log and publishes change events to Kafka.
- A streaming job reads micro-batches and deduplicates to the latest change per key.
- MERGE applies inserts, updates and deletes to the Delta table.
- Analysts query the Delta table.
The design follows the CDC system design case study; read it before building.
Progress is saved in this browser only. No account needed.