Skip to content
Good engineers know 4 min read · checkpointLocation 2 practice problems ↓

Structured Streaming basics

Streams as an unbounded table: sources, sinks, triggers, output modes and checkpoints.

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

Comfortable with these? Read on.

TL;DR Structured Streaming lets you write a stream with the same DataFrame code as a batch job. Spark treats the input as an unbounded table, runs your query on each new chunk of data (a micro-batch), and records progress in a checkpoint so a restarted query continues exactly where it stopped, with no lost or duplicated output when the sink supports it.

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.

PySpark · Kafka to a Delta table
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

PartOptionsNotes
SourceKafka, files in a folder (Auto Loader on Databricks), Delta tables, sockets/rate for testingMust be replayable (offsets) for exactly-once
SinkDelta / 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 testingDelta 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
TriggerprocessingTime="1 minute", availableNow=True, default (as fast as possible), realTime modes on some platformsavailableNow processes everything available and stops: a batch job with streaming bookkeeping
Output modeappend, update, completeSee below
CheckpointcheckpointLocationOne per query, on durable storage, never shared

Output modes

ModeWrites each batchWorks with
appendOnly new rows that will never change againMaps 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 →
updateOnly result rows that changed in this batchAggregations; sinks that can upsert (foreachBatch + MERGE)
completeThe whole result tableSmall aggregations only: it grows forever

Checkpoints and exactly-once

Step 1 · plan

The driverdriver: The process that runs your program, plans the work and sends tasks to executors. Learn more → asks the source what is new (Kafka offsets 1,000-1,450) and writes that range to the checkpoint's offsets log before processing. This is the write-ahead log.

Step 2 · process

Executorsexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more → read exactly that range, run the query and update state (stored in the checkpoint's state folder).

Step 3 · write

The sink writes the output. A Delta sink records the batch id with the commit, so writing batch 42 twice is detected and skipped.

Step 4 · commit

The batch id is written to the commits log. The batch is done.

Crash · restart

If the job dies between steps 1 and 4, the restarted query finds batch 42 in offsets but not in commits, and re-runs exactly the same offset range. Replayable source + deterministic processing + idempotent sink = exactly-once output.
Watch out: the checkpoint is tied to the query. Deleting it re-reads from 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:

PySpark · Upsert every batch
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.lastProgress reports 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

Lost or duplicated data. One checkpoint per query, kept as long as the query lives.

complete mode on a large aggregation

The whole result is rewritten every batch.

Non-idempotent foreachBatch

A re-run batch appends twice. Use MERGE or write with the batch id.

Aggregating without a watermark

State grows forever and eventually the job runs out of memory.

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 questions

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

Solve: Upsert: Merge New Records into a Table →

Keep going

Up next · lesson 13 of 24 · 4 min read
Watermarks and late data
Event time vs processing time, how watermarks bound state, and what happens to late rows.

Related lessons

Previous: The small file problem

Primary sources: Structured Streaming guide · Kafka integration