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).