4 · Advanced · lesson 18 of 20

Structured Streaming

The same DataFrame API — but the table keeps growing.

Structured Streaming treats a stream as a table that never stops appending. You write the same DataFrame code; Spark runs it incrementally on each micro-batch.

Python
events = (spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "broker:9092")
    .option("subscribe", "orders")
    .load()
    .selectExpr("CAST(value AS STRING) AS json"))

parsed = events.select(F.from_json("json", schema).alias("o")).select("o.*")

per_minute = (parsed
    .withWatermark("ts", "10 minutes")
    .groupBy(F.window("ts", "1 minute"), "country")
    .agg(F.sum("amount").alias("revenue")))

query = (per_minute.writeStream
    .outputMode("update")
    .format("console")
    .trigger(processingTime="30 seconds")
    .start())

query.awaitTermination()
TIP
withWatermark tells Spark how late events can arrive, so it can drop old state and keep aggregations bounded.
Loading 3D scene…
Key takeaways
  • ✓Streams = unbounded DataFrames; same API, same optimizer.
  • ✓Watermarks bound state so aggregates don't grow forever.
  • ✓outputMode: append (new rows only), update (changed rows), complete (whole result).