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.

Advertisement

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.

One micro-batch, from trigger to commit: what is written where, and in what order1. Trigger firesdriver wakes up2. Ask sourceslatest offsets3. Write offsets/Nplanned range, before any work4. Run incremental plan as a jobmap stages, shuffle, stateful stageState store, version None store per shuffle partitionread N-1, write N5. Sink writes batch Nidempotent by batch id, or transactional6. Write commits/Nbatch N is doneCheckpoint directoryoffsets/ commits/ state/ sources/ metadataOn restart: offsets/N present but commits/N missing means batch N is re-run with exactly the same offsets.That replay is why sinks must tolerate seeing batch N twice.
The driver writes the planned offsets before doing any work and the commit marker after the sink succeeds. Everything between the two may be repeated.
  1. The trigger fires. With no trigger set, the next batch starts as soon as the previous one ends.
  2. The driver asks each source for its latest available offsets and decides the range for batch N, capped by options such as maxOffsetsPerTrigger.
  3. It writes that range to offsets/N in the checkpoint directory. This is a write-ahead log: the plan is durable before any data is read.
  4. 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.
  5. The sink writes the batch's output. File and Delta sinks record the batch id so a repeat is skipped; foreachBatch hands you the batch id so you can do the same.
  6. 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.

Advertisement

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.

StepWhat costs timeIllustrative cost
Wait for triggerup to the trigger interval0 to 30 s with a 30 s trigger
Fetch latest offsetsone metadata round trip per sourcetens of ms
Write offsets/Na small file on the checkpoint storetens to hundreds of ms on object storage
Schedule the jobone task per partition per stagegrows with spark.sql.shuffle.partitions
Shuffle and statenetwork shuffle, state reads and writesdepends on data and state size
State commitupload state changes to the checkpointtens of ms to seconds
Sink and commits/Noutput write plus another small filetens 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

TriggerBehaviourUse it for
default (unset)next micro-batch starts as soon as the previous endslowest micro-batch latency
processingTime="30 seconds"batches start on a fixed schedulepredictable cost and file sizes
availableNow=Trueprocess everything available in several bounded batches, then stopscheduled incremental jobs; replaces once
once=Trueone batch of everything, then stop; deprecatednothing new
continuous="1 second"experimental long-running tasks, at-least-once, map-like queries onlyrarely worth it now
Trigger.RealTime(...)real-time mode, added in 4.1 (trigger name per current Apache docs): long-running tasksmillisecond-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 stateOperators in the progress events; numRowsTotal should 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 foreachBatch function 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, availableNow on 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?

ConcernSpark Structured StreamingFlink
Latencyseconds in micro-batch; milliseconds in real-time mode for stateless queriesmilliseconds, including stateful jobs
Programming modelsame DataFrame and SQL code as batchDataStream and Table APIs, separate from your batch stack
Custom state and timerstransformWithState since 4.0mature keyed state, timers and process functions
Recoveryreplay the last batch from the checkpointrestore from the last distributed snapshot
Operationsone engine and cluster for batch and streaminga 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

  1. List every streaming query and its trigger. Move scheduled incremental jobs to availableNow.
  2. 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.
  3. Switch large-state queries to the RocksDB provider with changelog checkpointing, starting with a new checkpoint.
  4. Add a progress listener and alert on batch duration above the trigger interval and on input rate above processing rate.
  5. Confirm every foreachBatch sink is idempotent by batch id, and kill a running query once to prove it.
  6. Read a production checkpoint with the state data source and find your largest keys.
  7. If a stateless pipeline needs milliseconds, prototype real-time mode on your Spark version before adding a second engine.
Key takeaway: Structured Streaming runs your DataFrame query incrementally, one micro-batch at a time: write the planned offsets, run the job, update versioned state, write the sink, write the commit. That order is what gives replay-based exactly-once output, and its fixed per-batch costs explain why latency is measured in seconds. Size shuffle partitions before the first start, use RocksDB with changelog checkpointing for real state, make sinks idempotent by batch id, and watch batch duration against the trigger. Reach for real-time mode for stateless millisecond work and for Flink when stateful logic needs it.