Spark Streaming is two things that share a name. The original DStreams API, deprecated since Spark 3.4, cut a stream into RDDs every few seconds. Structured Streaming, the engine every new pipeline should use, treats a stream as a table that keeps growing and runs your DataFrame query incrementally over it. The second one is the subject here: what happens inside one unit of work, where the time goes, what is durable, and what breaks.
Several neighbouring pages go deeper on single parts, and this one links to them rather than repeating them: Structured Streaming with Kafka for offsets and delivery guarantees, the watermark architecture for event time and late data, and Structured Streaming state for the state store itself. API names below were checked against the current Spark documentation (4.2 at the time of writing); where a feature arrived in a specific release, the release is named.
From DStreams to an unbounded table
DStreams exposed the batches directly. You wrote reduceByKeyAndWindow over RDDs, the batch interval was fixed when the context started, windows were based on arrival time, and keeping state meant updateStateByKey or mapWithState with checkpointing you largely managed yourself.
Structured Streaming moved all of that into the engine. You write an ordinary DataFrame query. The planner turns it into an incremental plan: for each new slice of input, compute what changed in the result and emit it according to the output mode (append, update or complete). Event time is a column, watermarks are declared with withWatermark, and state lives in a versioned store that is checkpointed together with the source offsets. The result is the property that makes the engine useful: if your sources can be replayed and your sink is idempotent or transactional, a crash at any point produces the same output as a run without a crash.
The micro-batch lifecycle, step by step
By default the engine runs in micro-batches. Each batch follows the same six steps, and the order of the writes is what makes recovery work.
- The trigger fires. With no trigger set, the next batch starts as soon as the previous one ends.
- The driver asks each source for its latest available offsets and decides the range for batch N, capped by options such as
maxOffsetsPerTrigger. - It writes that range to
offsets/Nin the checkpoint directory. This is a write-ahead log: the plan is durable before any data is read. - The incremental plan runs as an ordinary Spark job. Stateless stages read and transform; a stateful operator sits after a shuffle and each of its tasks opens the state store for its partition at version N-1 and writes version N.
- The sink writes the batch's output. File and Delta sinks record the batch id so a repeat is skipped;
foreachBatchhands you the batch id so you can do the same. - The driver writes
commits/N. Only now is batch N finished.
On restart, the engine reads the last offsets and commits entries. If offsets/N exists without commits/N, batch N is re-run over exactly the same offset range against state version N-1. That is the whole exactly-once story: deterministic replay plus a sink that tolerates a repeat.
Where one batch spends its time
Seconds of latency come from the fixed costs of the lifecycle. The breakdown below is illustrative, for a small stateful query writing its checkpoint to object storage; your numbers will differ, and the point is which parts are fixed regardless of how little data arrived.
| Step | What costs time | Illustrative cost |
|---|---|---|
| Wait for trigger | up to the trigger interval | 0 to 30 s with a 30 s trigger |
| Fetch latest offsets | one metadata round trip per source | tens of ms |
| Write offsets/N | a small file on the checkpoint store | tens to hundreds of ms on object storage |
| Schedule the job | one task per partition per stage | grows with spark.sql.shuffle.partitions |
| Shuffle and state | network shuffle, state reads and writes | depends on data and state size |
| State commit | upload state changes to the checkpoint | tens of ms to seconds |
| Sink and commits/N | output write plus another small file | tens to hundreds of ms |
Worked example: a query left at the default of 200 shuffle partitions launches 200 stateful tasks every batch, each of which opens and commits a state store, even when the batch carries a few hundred rows. Setting the partition count to roughly the number of cores available for the stateful stage often cuts batch time sharply. Do it before the first start: the number of state partitions is fixed when the checkpoint is created, and changing spark.sql.shuffle.partitions later does not repartition existing state.
The other number to watch is the ratio between batch duration and trigger interval. If a 30-second trigger consistently produces 40-second batches, the query starts each batch immediately after the last, input piles up, and latency grows without bound. Add parallelism or do less per record; a shorter trigger will not help.
Trigger modes and what each one is for
| Trigger | Behaviour | Use it for |
|---|---|---|
| default (unset) | next micro-batch starts as soon as the previous ends | lowest micro-batch latency |
processingTime="30 seconds" | batches start on a fixed schedule | predictable cost and file sizes |
availableNow=True | process everything available in several bounded batches, then stop | scheduled incremental jobs; replaces once |
once=True | one batch of everything, then stop; deprecated | nothing new |
continuous="1 second" | experimental long-running tasks, at-least-once, map-like queries only | rarely worth it now |
Trigger.RealTime(...) | real-time mode, added in 4.1 (trigger name per current Apache docs): long-running tasks | millisecond-latency stateless pipelines |
availableNow deserves more use than it gets. Run hourly, it is a batch job with exactly-once incremental semantics: the checkpoint remembers where it stopped, so no custom high-water-mark logic.
Real-time mode: when micro-batches are too slow
Spark 4.1 added a real-time mode that keeps the Structured Streaming API and guarantees but changes execution. Instead of waking on an interval, the engine launches one long-running task per input partition; each task pulls records continuously and pushes them through the transformation to the sink as they arrive. The duration you pass is a checkpoint interval: how often progress is committed and a new long-running batch begins. It is not a latency target. The current documentation sets a minimum of five seconds through spark.sql.streaming.realTimeMode.minBatchDuration.
// Scala, as in the current Apache docs; vendor builds may name the trigger differently
import org.apache.spark.sql.streaming.Trigger
enriched.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("topic", "checkouts-enriched")
.option("checkpointLocation", "/chk/enrich-rtm")
.outputMode("update") // the only mode real-time mode accepts
.trigger(Trigger.RealTime("5 minutes")) // a checkpoint interval, not a latency target
.start()The limits in the 4.1 release are strict: stateless, single-stage queries written in Scala, the Kafka source, the Kafka and foreach sinks, and update output mode only. Unions and broadcast stream-static joins are supported; aggregations, deduplication and stream-stream joins are not. The master-branch documentation is moving faster than the release, so check the page for your exact version before planning around anything stateful. Because each task stays resident, plan cluster slots for every source partition at once rather than sharing slots across stages over time.
State: providers, checkpoints and reading it back
Every stateful operator keeps one state store per shuffle partition. The default provider holds state as Java objects in executor memory and writes delta files to the checkpoint each batch. It is fast for small state and painful for large state, because millions of live objects mean long garbage-collection pauses. The RocksDB provider keeps state in native memory and local disk and uploads changed files on commit; it is the sensible default for anything larger than a few hundred megabytes per executor.
With RocksDB, enable changelog checkpointing. Instead of uploading snapshot files at every commit, each batch uploads only a log of the changes and snapshots are taken in the background, which makes commit time depend on what changed rather than on total state size. Both settings are in the pipeline below.
Spark 4.0 added a state data source for when you need to see what is actually held. spark.read.format("state-metadata").load(path) lists the stateful operators, their partitions and the batch ids available; spark.read.format("statestore").load(path) returns keys and values as a DataFrame you can aggregate. The metadata is only written by queries running on 4.0 or later.
Custom state with transformWithState
Built-in aggregations, deduplication and joins manage their own state. For anything else, the older flatMapGroupsWithState gave you one state object per key and a single timeout. Spark 4.0 added its successor, transformWithState (transformWithStateInPandas in Python): a stateful processor can declare several named value, list and map states, register many timers per key, and set a time-to-live on state.
The sketch below raises an alert for a store that has gone quiet. Each input deletes the key's existing timer, records the time and registers a new timer ten minutes out; if no input arrives before it fires, handleExpiredTimer emits the alert. The method names and arguments follow the 4.0 guide.
class SilentStore(StatefulProcessor):
"""Emit an alert when a store has sent no checkout for 10 minutes of processing time."""
def init(self, handle):
self.handle = handle
self.last = handle.getValueState(
"last_seen", StructType([StructField("ms", LongType())]))
def handleInputRows(self, key, rows, timerValues):
now = timerValues.getCurrentProcessingTimeInMs()
for t in self.handle.listTimers(): # one live timer per key
self.handle.deleteTimer(t)
self.last.update((now,))
self.handle.registerTimer(now + 600_000)
yield pd.DataFrame()
def handleExpiredTimer(self, key, timerValues, expiredTimerInfo):
yield pd.DataFrame({"store_id": [key[0]], "silent_since_ms": [self.last.get()[0]]})
def close(self):
pass
alerts = (events.groupBy("store_id")
.transformWithStateInPandas(
statefulProcessor=SilentStore(),
outputStructType=StructType([StructField("store_id", StringType()),
StructField("silent_since_ms", LongType())]),
outputMode="Update",
timeMode="ProcessingTime"))
A worked pipeline, with the knobs that matter
The query below reads checkouts from Kafka, counts orders and revenue per store per minute of event time, and writes through foreachBatch. Every option in it is there for a reason covered above: a fixed shuffle partition count chosen up front, the RocksDB provider with changelog checkpointing, a per-batch cap so a restart after an outage does not attempt hours of backlog in one batch, a watermark so window state is dropped, and a listener so you can see batch durations.
from pyspark.sql import SparkSession, functions as F, types as T
from pyspark.sql.streaming import StreamingQueryListener
spark = (SparkSession.builder.appName("checkout-metrics")
.config("spark.sql.shuffle.partitions", "48") # fixed for the life of the checkpoint
.config("spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
.config("spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled", "true")
.getOrCreate())
schema = T.StructType([
T.StructField("store_id", T.StringType()),
T.StructField("amount", T.DoubleType()),
T.StructField("ts", T.TimestampType()),
])
events = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "checkouts")
.option("maxOffsetsPerTrigger", "200000") # bounds the catch-up batch
.load()
.select(F.from_json(F.col("value").cast("string"), schema).alias("e"))
.select("e.*"))
per_minute = (events
.withWatermark("ts", "2 minutes")
.groupBy(F.window("ts", "1 minute"), "store_id")
.agg(F.count("*").alias("orders"), F.sum("amount").alias("revenue")))
def upsert(batch_df, batch_id):
# Called once per micro-batch; may be called again with the same batch_id after a failure.
(batch_df.withColumn("batch_id", F.lit(batch_id))
.write.mode("append").format("delta").save("/lake/checkout_minutes_staging"))
class Progress(StreamingQueryListener):
def onQueryStarted(self, event): pass
def onQueryProgress(self, event):
p = event.progress
print(p.batchId, p.numInputRows, p.inputRowsPerSecond,
p.processedRowsPerSecond, p.durationMs.get("triggerExecution"))
def onQueryIdle(self, event): pass
def onQueryTerminated(self, event): print("terminated", event.exception)
spark.streams.addListener(Progress())
query = (per_minute.writeStream
.outputMode("update")
.foreachBatch(upsert)
.option("checkpointLocation", "/chk/checkout-metrics")
.trigger(processingTime="30 seconds")
.start())Two details are easy to miss. The staging write tags rows with the batch id, so a downstream merge can discard a repeated batch; a plain append would duplicate rows after any retry. And the listener prints triggerExecution, the wall time of the batch. Alert when it stays above the trigger interval, and when inputRowsPerSecond exceeds processedRowsPerSecond for several minutes: those two conditions mean the query is falling behind.
Failure modes seen in production
- State that never shrinks. An aggregation or deduplication without a watermark keeps every key forever. Look at
stateOperatorsin the progress events;numRowsTotalshould plateau. - A checkpoint that will not restart. Changing the stateful parts of the query, the shuffle partition count or the source identity makes the old checkpoint incompatible. Plan such changes as a new query with a new checkpoint and a backfill.
- Duplicate output after a crash. A
foreachBatchfunction that writes to a database without using the batch id produces duplicates on replay. - Small files. A file sink with a short trigger writes at least one file per partition per batch. Use a longer trigger,
availableNowon a schedule, or a table format with compaction. - Driver loss. The driver plans every batch. Run it under a supervisor that restarts it, which is safe because the checkpoint holds the progress.
Spark or Flink?
| Concern | Spark Structured Streaming | Flink |
|---|---|---|
| Latency | seconds in micro-batch; milliseconds in real-time mode for stateless queries | milliseconds, including stateful jobs |
| Programming model | same DataFrame and SQL code as batch | DataStream and Table APIs, separate from your batch stack |
| Custom state and timers | transformWithState since 4.0 | mature keyed state, timers and process functions |
| Recovery | replay the last batch from the checkpoint | restore from the last distributed snapshot |
| Operations | one engine and cluster for batch and streaming | a separate cluster and expertise |
Choose Spark when your team already runs Spark, your latency budget is seconds, and sharing code between batch and streaming matters. Choose Flink when you need low latency for stateful logic, very large keyed state, or rich event-time processing; Flink state architecture shows what you get in exchange for running a second engine.
What to do next
- List every streaming query and its trigger. Move scheduled incremental jobs to
availableNow. - For each stateful query, record the shuffle partition count fixed in its checkpoint and compare it with the cores you actually give the stateful stage.
- Switch large-state queries to the RocksDB provider with changelog checkpointing, starting with a new checkpoint.
- Add a progress listener and alert on batch duration above the trigger interval and on input rate above processing rate.
- Confirm every
foreachBatchsink is idempotent by batch id, and kill a running query once to prove it. - Read a production checkpoint with the state data source and find your largest keys.
- If a stateless pipeline needs milliseconds, prototype real-time mode on your Spark version before adding a second engine.