1 December 2024 · 8 min · Kafka · Spark · Streaming
Building real-time data pipelines with Kafka and Spark
How I design and run streaming pipelines with Apache Kafka and Spark Structured Streaming for production workloads.
Why real time
Whether you're following IoT sensors, processing payments or watching how people use a product, being able to act on data the moment it arrives is a real advantage.
This post covers the architecture and patterns I've used to build production streaming pipelines with Apache Kafka and Spark Structured Streaming.
The shape of the pipeline
- ✦Ingestion — Kafka is the central hub that receives events from every source
- ✦Processing — Spark Structured Streaming cleans, joins and enriches the data
- ✦Storage — Delta Lake gives reliable storage with ACID guarantees
Exactly-once processing
The hardest part of streaming is making sure each event is counted once, not zero times and not twice. Checkpoints plus a transactional sink get you there:
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("subscribe", "events")
.option("startingOffsets", "earliest")
.load()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/events")
.trigger(processingTime="10 seconds")
.start("/data/events")Late data
Watermarks tell Spark how long to wait for events that arrive late:
df.withWatermark("event_time", "1 hour")
.groupBy(window("event_time", "5 minutes"))
.agg(count("*").alias("event_count"))What to watch
- ✦Latency — time from an event arriving to it showing up in the output
- ✦Throughput — events processed per second
- ✦Backpressure — how deep the queue in each Kafka partition is getting
Wrapping up
Reliable streaming comes down to error handling, state management and monitoring. Kafka and Spark together give you a solid base for all three.
Working on something like this?
Get in touch →