ProjectsProject 5 of 8
Advanced project · Project 5 of 8
Kafka → Spark → Delta Lake Streaming Pipeline
An application emits user events to Kafka. Build a streaming pipeline that lands them in Delta Lake within a minute, deduplicates replays, and produces per-minute aggregates that tolerate late events.
Requirements
- Consume events from a Kafka topic
- Write raw events to a Delta table with checkpointing
- Deduplicate events by event id within a watermark
- Compute per-minute aggregates with event-time windows
- Recover from restarts without losing or double-counting data
Technology stack
Kafka, Spark Structured Streaming, Delta Lake, Docker Compose for local Kafka
Dataset
Generate synthetic events with a small producer script.
Business context
Teams increasingly need minute-level data for product monitoring. This project covers the hardest parts of streaming in a small setting: state, late data, duplicates and recovery.
Architecture
- A producer sends keyed events to a Kafka topic.
- Spark Structured Streaming reads micro-batches and parses JSON.
- Raw events are appended to a Delta table; progress is tracked in a checkpoint.
- Deduplicated, windowed aggregates are written to a second Delta table.
Keep the checkpoint location stable: deleting it makes Spark treat the stream as new.
Progress is saved in this browser only. No account needed.