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.

Three parts: source, query, sink.
Figure 6.1 — Source, query and sink.

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.