Lesson 19 / 25
The Structured Streaming Model
Write streaming queries with the same DataFrame API and understand triggers and output modes.
Same API, unbounded input
Structured Streaming lets you write a streaming job as if it were a batch job on a table that new rows keep being appended to. You read with spark.readStream from a source (Kafka, files, a test rate source), apply ordinary DataFrame operations and write with writeStream to a sink (Kafka, files/Delta, a console or memory table for testing). A trigger controls how often micro-batches run (processingTime="10 seconds", or availableNow=True to process what is there and stop). The output mode says what is written each time: append (only new final rows), update (rows that changed) or complete (the whole result table, for small aggregations). A checkpoint location stores progress and state so a restarted query resumes without losing or duplicating data (with replayable sources and idempotent or transactional sinks).
A table that keeps growing
Structured Streaming treats a live stream as an unbounded table and runs your query incrementally.
A tiny streaming aggregation (illustrative)
Uses the built-in rate source (a test generator) and an in-memory sink. I did not capture this run, so no output is shown; counts keep growing as batches arrive. For real jobs you would write to Kafka or Delta with a checkpointLocation.
import time
q = (spark.readStream.format("rate").option("rowsPerSecond", 20).load()
.withColumn("bucket", F.col("value") % 3)
.groupBy("bucket").count()
.writeStream.outputMode("complete").format("memory").queryName("rates")
.trigger(processingTime="1 second").start())
time.sleep(5)
spark.sql("SELECT bucket, count FROM rates ORDER BY bucket").show()
q.stop()Always set a checkpoint in production
Without a checkpoint a restarted query starts over (or fails to resume correctly). Give each query its own durable checkpoint folder, and do not delete or share it casually.
Quick check: What does a checkpoint location store?
- Only the final output files
- Your passwords
- The Spark binaries
- Query progress and state so a restart can resume
Answer
Query progress and state so a restart can resume — Offsets and state in the checkpoint let a restarted query continue where it stopped.