LakeFlow · architecture
- Postgres
- Debezium
- Kafka
- Spark
- Bronze
- Silver
- Gold
- Quarantine
Airflow · nightly backfill
- CDC events replayed
- 510,663
- end-to-end to Gold (2 vCPU)
- 29.5 s
- events per second
- ~17.3K/s
- rows quarantined with reasons
- 5,239
Streams every insert, update and delete out of PostgreSQL through Debezium and Kafka into Spark Structured Streaming, landing a medallion Parquet lakehouse with SCD Type 2 history and a data-quality quarantine.
Problem
Operational databases change thousands of times a second, but analytics usually sees a nightly snapshot. Deletes vanish, history is lost, and replaying a failed run duplicates rows. The goal was a pipeline where the warehouse is a faithful, replayable, historised mirror of the OLTP system — in seconds, not hours.
Approach
- 01
Postgres runs with logical WAL and REPLICA IDENTITY FULL so Debezium emits before/after images; Kafka (KRaft) carries one topic per table.
- 02
Spark Structured Streaming appends the raw envelope to an immutable, checkpointed Bronze layer partitioned by table and ingest date — exactly-once at the sink.
- 03
foreachBatch upserts derive Silver with latest-per-key on the WAL commit order (ts_ms, lsn), honour deletes, and version customers as SCD Type 2. Rows failing declarative rules land in a quarantine table with the rule names.
- 04
Gold marts (daily revenue, product ranks, customer LTV) are Spark SQL; one set of transforms powers both the stream and the nightly Airflow backfill, so a single pytest suite covers both paths.
Results
- CDC events replayed
- 510,663
- end-to-end to Gold (2 vCPU)
- 29.5 s
- events per second
- ~17.3K/s
- rows quarantined with reasons
- 5,239
- Created
- 3 Sep 2026
- Commits
- 2
- Default branch
- main
- License
- unlicensed
Stack
- PySpark
- Structured Streaming
- Kafka
- Debezium
- PostgreSQL
- Parquet
- Airflow
- Docker Compose
- pytest