← Blog

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 →

Next post

RAG pipelines for business LLM applications

→