Structured Streaming is Spark's engine for continuous data. You write the same DataFrame or SQL query you would write for a table, declare a source that keeps growing and a sink that keeps receiving, and Spark runs that query incrementally: every trigger it reads only the new input, updates whatever state the query needs, and writes the changed results. Fault tolerance is not bolted on afterwards; it is the shape of the loop. Each batch's input range is logged before the work starts and its completion is logged after the sink write, so a restarted query knows precisely what to redo.

This page takes the engine apart: the micro-batch loop and checkpoint layout, event time and watermarks, the state store, the exactly-once contract between sources and sinks, triggers, and the operational signals that tell you a query is healthy. An earlier version of this page made several claims that are wrong for Structured Streaming (rate-based backpressure, adaptive query execution in streaming, a built-in side output for late events); they are corrected below.

One micro-batch, from offsets to commitSourcesKafka, files, Delta1. Plan batch Npick end offsets per sourceoffsets/Nwrite-ahead log entrywrite first2. Incremental plansame optimizer as batch3. Stateful operatorswindows, joins, dedupstate/ storeversioned per batchread + writeWatermarkmax event time - delayevict4. Sink writeidempotent by batch idcommits/Nbatch N is donethenRestart: replay batch N if offsets/N exists without commits/N, with the same offsetsand the same state version, so a sink that dedupes on batch id gives exactly-once output.
Each micro-batch logs its offsets before running and its completion after the sink write; the pair is what makes recovery deterministic.

The model: an unbounded table and an incremental plan

The programming model treats a stream as an unbounded input table to which rows are appended. A query over it defines a result table, and the output mode says what to emit when the result changes. append emits only rows that will never change again, update emits rows that changed in this batch, and complete rewrites the whole result each time, which only makes sense for small aggregates.

The planner turns the query into an incremental plan. Stateless operators such as select, filter and stream-static joins simply process the new rows. Stateful operators, including aggregations, stream-stream joins, dropDuplicates and arbitrary stateful processing, keep partial results in a state store between batches. Some batch features do not exist in streaming: global sorts, limit in most positions, and adaptive query execution, which Spark disables for streaming queries because a stateful plan's partitioning must not change between batches.

Joins show how the model shapes cost. A stream-static join, enriching clicks with a campaign dimension table, is stateless: each batch joins its new rows against the static side, and a Delta dimension table is re-read per batch so updates appear without a restart. A stream-stream join is stateful, because a click may arrive before or after the impression it matches. Both sides are buffered in state until a match can no longer occur, which only becomes knowable when both inputs carry watermarks and the join condition includes an event-time bound such as click_time BETWEEN imp_time AND imp_time + interval 1 hour. Leave out either and the buffered rows are never evicted, so state grows with total traffic rather than with the time window.

The model: an unbounded table and an incremental plan

The programming model treats a stream as an unbounded input table to which rows are appended. A query over it defines a result table, and the output mode says what to emit when the result changes. append emits only rows that will never change again, update emits rows that changed in this batch, and complete rewrites the whole result each time, which only makes sense for small aggregates.

The planner turns the query into an incremental plan. Stateless operators such as select, filter and stream-static joins simply process the new rows. Stateful operators, including aggregations, stream-stream joins, dropDuplicates and arbitrary stateful processing, keep partial results in a state store between batches. Some batch features do not exist in streaming: global sorts, limit in most positions, and adaptive query execution, which Spark disables for streaming queries because a stateful plan's partitioning must not change between batches.

Joins show how the model shapes cost. A stream-static join, enriching clicks with a campaign dimension table, is stateless: each batch joins its new rows against the static side, and a Delta dimension table is re-read per batch so updates appear without a restart. A stream-stream join is stateful, because a click may arrive before or after the impression it matches. Both sides are buffered in state until a match can no longer occur, which only becomes knowable when both inputs carry watermarks and the join condition includes an event-time bound such as click_time BETWEEN imp_time AND imp_time + interval 1 hour. Leave out either and the buffered rows are never evicted, so state grows with total traffic rather than with the time window.

The micro-batch loop and the checkpoint

The default execution mode is micro-batch. A driver-side loop runs one Spark job per batch, and each batch follows the same protocol.

loop:
    wait for trigger                                  # interval, AvailableNow, or immediately
    if commits/ lacks the last batch in offsets/:
        N = last planned batch; end = offsets[N]      # crash recovery: redo exactly this batch
    else:
        N = last + 1
        end = {src: src.latestOffset(limits) for src in sources}   # e.g. maxOffsetsPerTrigger
        if no new data and no time-based work: continue
        write offsets/N = (end, watermark, conf)      # write-ahead log
    run job: read (start, end] -> operators -> state store version N -> sink.addBatch(N)
    write commits/N                                   # batch N done
    advance watermark from max event time seen

The checkpoint directory, which you name with checkpointLocation, contains metadata (the query id), offsets/, commits/, sources/ (for example the file source's seen-files log) and state/, with one subdirectory per stateful operator and partition. It must live on durable storage that every restart can reach, typically object storage or HDFS. The checkpoint is also a contract: the number of state partitions is fixed by spark.sql.shuffle.partitions at the first run and cannot change later, and changing stateful operators, grouping keys or the schema of state can make an existing checkpoint unusable.

Event time and watermarks

Processing time is when Spark sees a row; event time is a column in the row saying when it happened. Windows should use event time, because data arrives late and out of order. The engine therefore needs to know when a window is finished, and that is what a watermark declares. With withWatermark("event_time", "10 minutes"), at the end of each batch Spark computes the maximum event time it has seen and subtracts ten minutes. That value becomes the watermark for the next batch. If a query has several watermarked inputs, the global watermark defaults to the minimum of them (spark.sql.streaming.multipleWatermarkPolicy).

The watermark does two jobs. In append mode it decides when a window is final and may be emitted, so output lags by at least the delay. In every mode it lets the engine evict state for windows, join rows and dedup keys that can no longer change. The guarantee is one-sided: rows later than the delay may be dropped, and rows within it are never dropped. Dropped rows are not routed anywhere; Spark has no late-data side output. You can count them through numRowsDroppedByWatermark in query progress, and if you must keep them, write the raw stream to a separate sink as well. Without a watermark, aggregation and dedup state grows forever. The watermark article works through the edge cases.

The state store

State lives in a versioned key-value store per operator partition, co-located with the task that owns that partition. Each batch reads version N-1, applies updates and commits version N into the checkpoint. Open source Spark ships two providers. The default HDFSBackedStateStoreProvider keeps state in executor JVM memory and writes delta files plus periodic snapshots; it is fast for small state but turns large state into garbage-collection pauses. The RocksDBStateStoreProvider keeps state in native RocksDB on local disk and uploads files to the checkpoint, which handles state far larger than the heap. Newer releases add changelog checkpointing for RocksDB, which uploads only the changes made in a batch instead of fresh snapshot files, cutting commit latency.

spark.conf.set("spark.sql.streaming.stateStore.providerClass",
    "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
spark.conf.set("spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled", "true")

Set the provider before the first run; it is recorded with the checkpoint. For custom logic, Spark 4.0 added transformWithState (transformWithStateInPandas in Python), the successor to flatMapGroupsWithState, with typed value, list and map state, timers and TTL eviction. The state store article covers sizing and the state data source, which lets you query a checkpoint's state as a table when debugging.

The exactly-once contract

End-to-end exactly-once needs three properties together. The source must be replayable from a recorded offset, as Kafka, Kinesis, files and Delta tables are; a socket source is not. The engine must be deterministic over a given offset range and state version, which the offset log and versioned state provide. And the sink must be idempotent or transactional per batch id, so that redoing batch N after a crash does not duplicate its output.

SinkGuaranteeHow
Delta LakeExactly-onceCommits record query id and batch id; a replayed batch is skipped
File sinkExactly-onceA _spark_metadata log lists committed files; readers use it
KafkaAt-least-onceA replay can write duplicates; dedupe downstream by key
foreachBatchYour code decidesUse batchId to make writes idempotent

foreachBatch hands you each micro-batch as a normal DataFrame plus its id. It is the escape hatch for MERGE upserts, JDBC, and writing to several tables, and it is where most duplicate bugs are born, because the function runs again with the same id after a failure. The foreachBatch article shows idempotent patterns.

Triggers and rate limits

TriggerBehaviourUse it for
Default (unspecified)Next batch starts when the last endsLowest micro-batch latency
ProcessingTime("30 seconds")Fixed interval; a slow batch delays the nextPredictable cost, fewer small files
AvailableNowProcesses all data available at start in several batches, then stopsScheduled incremental jobs
OnceOne batch, then stop; deprecated in favour of AvailableNowLegacy jobs only
Continuous("1 second")Experimental long-running tasks, map-like queries onlyRarely; check its limits first

Rate limiting is explicit, not adaptive. Structured Streaming does not implement the DStreams-era backpressure setting; you cap batch size with source options such as maxOffsetsPerTrigger for Kafka or maxFilesPerTrigger for files. More on choosing in the triggers article.

Worked example: hourly clicks from Kafka to Delta

A clickstream topic must become hourly clicks per campaign in a Delta table, tolerating events up to fifteen minutes late.

from pyspark.sql import functions as F

clicks = (spark.readStream.format("kafka")
    .option("kafka.bootstrap.servers", "broker:9092")
    .option("subscribe", "clicks")
    .option("startingOffsets", "earliest")
    .option("maxOffsetsPerTrigger", 200000)
    .load()
    .select(F.from_json(F.col("value").cast("string"),
                        "campaign_id STRING, user_id STRING, event_time TIMESTAMP").alias("e"))
    .select("e.*"))

hourly = (clicks
    .withWatermark("event_time", "15 minutes")
    .groupBy(F.window("event_time", "1 hour"), "campaign_id")
    .count())

query = (hourly.writeStream.format("delta")
    .outputMode("append")
    .option("checkpointLocation", "s3://lake/_chk/hourly_clicks")
    .trigger(processingTime="1 minute")
    .start("s3://lake/gold/hourly_clicks"))

At 10:20 the newest event seen is 10:19, so the watermark is 10:04. The 09:00-10:00 window ended before 10:04 and is emitted and evicted from state; the 10:00-11:00 window stays open. An event stamped 09:58 arriving now is behind the watermark and is dropped. An event stamped 10:07 updates the open window. Now the driver dies after writing offsets/42 and part of the Delta write, but before the Delta commit and commits/42. On restart Spark finds batch 42 planned but not committed, re-reads exactly the same Kafka offsets, loads state version 41 and recomputes. Delta has no commit for batch 42 from this query, so it accepts the new one; had the commit landed before the crash, Delta would skip the replay. Either way each window appears once.

Failure modes

  • State explosion. A missing watermark, or dedup on a key without event time, grows state until executors fail. Watch numRowsTotal per state operator.
  • Data loss on Kafka retention. If the query falls behind retention, offsets vanish. failOnDataLoss (default true) stops the query; setting it false silently skips data.
  • Incompatible restart. Changing grouping keys, aggregations or shuffle partitions against an existing checkpoint fails or misbehaves. Plan migrations as a new checkpoint plus backfill.
  • Duplicates from foreachBatch. Non-idempotent writes repeat on retry. Key them on batch id.
  • Small files. Short triggers into file or Delta sinks create many tiny files. Lengthen the trigger and schedule compaction.
  • Falling behind. Batch duration above the trigger interval means input outpaces processing; latency rises until you add capacity or partitions.

Operating a streaming query

Every batch produces a progress report, available from query.lastProgress and from a StreamingQueryListener. Export it. The health rule is simple: processedRowsPerSecond must stay above inputRowsPerSecond, and durationMs.triggerExecution must stay below the trigger interval. Also watch the event-time watermark against wall-clock time, state rows and memory per operator, and Kafka consumer lag. For Kafka inputs, size partitions so each executor core has work; the Kafka source article covers offsets and partition mapping.

The trade-offs are mostly about latency against cost and correctness. A longer watermark accepts later data but holds more state and delays append output. A shorter trigger lowers latency but adds per-batch overhead and small files. RocksDB state handles scale at some per-access cost compared with in-heap state. Update mode gives fresher results but forces the sink to handle upserts.

What to do next

  1. Give every query a dedicated, durable checkpoint location and never share it.
  2. Use event time and a watermark for every aggregation, join and dedup.
  3. Pick output mode by what the sink can handle: append for immutable files, update with upserts.
  4. Use a replayable source and an idempotent sink; key foreachBatch writes on batch id.
  5. Set maxOffsetsPerTrigger and a trigger interval deliberately.
  6. Choose the state provider and shuffle partitions before the first run.
  7. Alert on input versus processed rate, batch duration, state rows and Kafka lag.
  8. Rehearse a restart and a query change in staging before production.
Key takeaway: Structured Streaming runs a batch query incrementally, logging each batch's offsets before work and its commit after the sink write. Exactly-once output follows only when a replayable source, versioned state and an idempotent sink line up. Watermarks bound both lateness and state, rate limits are explicit options, and the checkpoint is a contract you must plan changes around.