Structured Streaming basics
Streams as an unbounded table: sources, sinks, triggers, output modes and checkpoints.
On this page
You will learn
- How Structured Streaming treats a stream as a table that keeps growing
- Sources, sinks, triggers and output modes
- What a checkpoint stores and why it gives exactly-once results
- How to monitor a stream and keep it healthy
Read first
- Transformations vs actions · 3 min read
- The small file problem · 6 min read
Comfortable with these? Read on.
A stream is a table that keeps growing
Every new event is a new row appended to an input table that never ends. Your query is defined on that table, and Spark updates the result incrementally: it processes only the new rows each time and keeps whatever state it needs (running counts, open windows) between batches.
raw = (spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "broker:9092") .option("subscribe", "orders") .option("startingOffsets", "earliest") .load()) orders = (raw .select(F.from_json(F.col("value").cast("string"), order_schema).alias("o")) .select("o.*")) query = (orders.writeStream .format("delta") .outputMode("append") .option("checkpointLocation", "s3://chk/orders/") .trigger(processingTime="1 minute") .toTable("bronze.orders"))
The parts
| Part | Options | Notes |
|---|---|---|
| Source | Kafka, files in a folder (Auto Loader on Databricks), Delta tables, sockets/rate for testing | Must be replayable (offsets) for exactly-once |
| Sink | Delta / ParquetParquet: A columnar file format: values of each column are stored together, so queries read only the columns they need. Learn more → files, Kafka, foreachBatch, console and memory for testing | Delta and file sinks are idempotentidempotent: Safe to run again: running a step twice gives the same result as running it once. Learn more →; Kafka is at-least-once |
| Trigger | processingTime="1 minute", availableNow=True, default (as fast as possible), realTime modes on some platforms | availableNow processes everything available and stops: a batch job with streaming bookkeeping |
| Output mode | append, update, complete | See below |
| Checkpoint | checkpointLocation | One per query, on durable storage, never shared |
Output modes
| Mode | Writes each batch | Works with |
|---|---|---|
| append | Only new rows that will never change again | Maps and filters; aggregations only with a watermarkwatermark: In streaming, how late data may arrive. Spark uses it to know when old state can be dropped. Learn more → |
| update | Only result rows that changed in this batch | Aggregations; sinks that can upsert (foreachBatch + MERGE) |
| complete | The whole result table | Small aggregations only: it grows forever |
Checkpoints and exactly-once
Step 1 · plan
Step 2 · process
Step 3 · write
Step 4 · commit
Crash · restart
startingOffsets (duplicates or gaps), and some query changes, such as changing aggregation keys, are incompatible with an existing checkpoint. Plan changes as a new query with a new checkpoint and a deliberate starting point.foreachBatch: batch code per micro-batch
For sinks Spark does not support, or to MERGE into a table, foreachBatch hands you each micro-batch as a normal DataFrame:
def upsert(batch_df, batch_id): (DeltaTable.forName(spark, "silver.customers").alias("t") .merge(batch_df.alias("s"), "t.id = s.id") .whenMatchedUpdateAll().whenNotMatchedInsertAll() .execute()) updates.writeStream.foreachBatch(upsert) \ .option("checkpointLocation", "s3://chk/customers/").start()
Make the function idempotent (MERGE is; a plain append is not), because a batch can be re-run after a failure.
Keeping a stream healthy
query.lastProgressreports input rows per second, processed rows per second, batch duration and state size. If processing is slower than input, the stream falls behind.- Short triggers write many small files; compact the target table or use a longer trigger.
- State grows without bound for aggregations and dropDuplicates unless a watermark limits it.
- Kafka retention must be longer than your longest possible outage, or offsets you need are deleted.
Common mistakes
Sharing or deleting a checkpoint
complete mode on a large aggregation
Non-idempotent foreachBatch
Aggregating without a watermark
Key takeaways
- A stream is an unbounded table processed in micro-batches.
- Checkpoints store offsets, commits and state; never share or delete them casually.
- Exactly-once needs a replayable source, the checkpoint and an idempotent sink.
- availableNow turns a stream into an incremental batch job.
Check yourself
3 questions1. What does the offsets log in a checkpoint record?
Show the answer
The input range of each batch, written before processing. It is a write-ahead log so a failed batch can be re-run on exactly the same input.
2. Which output mode writes only rows that changed in this batch?
Show the answer
update. update emits changed result rows.
3. What does trigger(availableNow=True) do?
Show the answer
Processes all available data, then stops. It is a batch-style run that still uses the checkpoint to process only new data.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 13 of 24 · 4 min readWatermarks and late data
Event time vs processing time, how watermarks bound state, and what happens to late rows.
Related lessons
Streaming into Delta tablesSpark internals · 6 min read
The small file problem
Previous: The small file problem
Primary sources: Structured Streaming guide · Kafka integration