Menu

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.

  • Advanced
  • Kafka · Spark Structured Streaming · Delta Lake · Docker Compose for local Kafka
  • 2 min read
  • Updated Oct 2026

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

  1. A producer sends keyed events to a Kafka topic.
  2. Spark Structured Streaming reads micro-batches and parses JSON.
  3. Raw events are appended to a Delta table; progress is tracked in a checkpoint.
  4. Deduplicated, windowed aggregates are written to a second Delta table.
Kafka offsets and output commits are tracked together in the checkpoint, so restarts resume cleanly.

Keep the checkpoint location stable: deleting it makes Spark treat the stream as new.

By Data Career Hub Editorial · Last reviewed Oct 2026

Progress is saved in this browser only. No account needed.

Search
Filter by type