skip to content

Structured Streaming

Spark's streaming model treats an unbounded stream as a table that keeps growing, processed incrementally with triggers, watermarks and checkpoints. Interviewers probe output modes and late-data handling because that is where 'it worked as a batch job' quietly stops being true.

on this pageshow

explore

questions

6

In Spark Structured Streaming, what does each output mode — append, update, complete — write?

level: middleimportance: must knowfreq 74%

answer

  1. three ways to publish the result table
  2. one of them never revises a row
  3. one rewrites everything every trigger
  4. the middle one demands upserts downstream
  5. append plus aggregation needs a watermark

basics

~10 s

Append writes only rows that are final and never revises them. Update writes only rows whose value changed in this trigger. Complete rewrites the entire result table every trigger and requires an aggregation.

solid answer

~50 s

`outputMode` decides which part of the conceptual result table Spark pushes to the sink on each trigger. **Append** (the default) emits only rows that are new and will never change again; for a stateless `select`/`filter` that is every row immediately, but for a streaming aggregation Spark needs a watermark to know when a group is finished, and rejects the query at start without one. **Update** emits every row whose value changed in this trigger, so an aggregation streams out successive revisions of the same key — the sink has to upsert by key or the consumer has to take the last value. **Complete** rewrites the whole result table every trigger; it requires an aggregation and keeps that aggregation state forever, so it only suits small key spaces. Sinks constrain the choice too: the file sink accepts append only.

code

python · 14 lines
python
agg = (spark.readStream.format("kafka")
       .option("subscribe", "clicks").load()
       .selectExpr("CAST(value AS STRING) AS body", "timestamp AS ts")
       .withWatermark("ts", "10 minutes")
       .groupBy(window("ts", "5 minutes"), "body")
       .count())

# final rows only, gated by the watermark
agg.writeStream.outputMode("append").format("parquet") \
   .option("path", "/out").option("checkpointLocation", "/ckpt/a").start()

# revisions as they happen; sink must upsert by key
agg.writeStream.outputMode("update").foreachBatch(upsert_by_key) \
   .option("checkpointLocation", "/ckpt/u").start()

go deeper

for a junior

Learn the three names and one sentence each: append adds final rows, update re-sends changed rows, complete rewrites everything. Know that append is the default and that the file sink accepts only append.

for a middle

Be ready to explain why an aggregation in append mode needs a watermark, and what each mode demands of the sink — that update-mode output is a stream of revisions per key and needs an upsert on the other end.

for a senior

Show that you pick a mode from the consumer's requirements and the latency budget, not by habit. Expect to talk about append's watermark-gated latency, complete mode's unbounded state, and how you would implement update semantics through foreachBatch.

for a principal

Own the mode as an interface decision across teams: append gives downstream immutable, replayable facts at higher latency; update gives speed but pushes idempotency onto every consumer. Be able to justify which contract your platform standardises on.

## The result-table model Structured Streaming asks you to think of the stream as a table that keeps growing. Each trigger appends the newly arrived records to a conceptual *input table*, Spark re-evaluates your query against the whole input table, and the answer is a conceptual *result table*. Spark does not literally materialise either one — the plan is executed incrementally, and only the state needed to keep the answer correct is retained. The **output mode** is the contract for which slice of that result table gets handed to the sink each trigger. That framing explains why the three modes exist. For a stateless query the result table only ever grows, so "what is new" is unambiguous. For an aggregation, an existing row can be *revised* when a later record lands in its group, and the sink now needs to be told whether it will receive revisions, only final values, or a full refresh. ## Append `outputMode("append")` is the default. It emits only rows added since the previous trigger **that are guaranteed never to change again**. For `select`, `filter`, `withColumn`, `explode` and other stateless transformations every output row qualifies the moment it is produced, so append behaves the way an ordinary batch write does. For a streaming aggregation the guarantee is the hard part. `groupBy(window($"ts", "10 minutes")).count()` cannot finalise a window's count until Spark believes no further records for that window will arrive — which is exactly what `withWatermark` declares. Write a streaming aggregation in append mode with no watermark and Spark refuses to start the query, reporting that append output mode is not supported for streaming aggregations without a watermark. With a watermark, the window's row is buffered in state and emitted once — after the watermark passes the window's end. The consequence is latency: a 10-minute window with a 5-minute watermark delay produces its first output roughly 15 minutes of event time after the window opened. Stream–stream outer joins have the same shape: the NULL-padded row can only be emitted once the watermark proves no match will arrive, so those joins require a watermark and a time-range condition. ## Update `outputMode("update")` emits every result row whose value changed during this trigger, and nothing else. Rows that were untouched are not re-sent. For an aggregation this means the same key is written repeatedly with a growing count — a partial answer, then a better one, then the final one. That is the lowest-latency mode for aggregations: you see a window's running count seconds after it opens rather than waiting for the watermark. The cost lands on the sink. Because the same key arrives many times, the sink must be an upsert (a key-value store, a `MERGE` in a table sink, a `foreachBatch` that writes by primary key) or the downstream reader must be tolerant of superseded values. Appending update-mode output to a plain file sink would give you every intermediate count as a separate row. A watermark is optional in update mode, but omitting it means aggregation state is never evicted and grows for the life of the query. ## Complete `outputMode("complete")` rewrites the entire result table on every trigger. It is only legal when the query contains an aggregation — there is no meaningful "whole table" for a projection over an unbounded stream. Because Spark must be able to re-emit every group forever, it retains all aggregation state regardless of any watermark you set; a watermark does not evict state in complete mode. That makes complete mode a small-cardinality tool: totals by country, counts by status code, a dashboard-shaped summary. Use it for per-user or per-session keys and state grows without bound and every trigger rewrites an ever-larger output. ## Sink compatibility The mode you want is not always available. The file sink supports append only, because files already committed cannot be rewritten in place. The Kafka sink accepts all three, though complete mode republishes the full table into the topic each trigger. The `console` and `memory` sinks accept all three and are debugging tools. `foreach`/`foreachBatch` accept all three and push the semantics onto your code — `foreachBatch` hands you a normal batch DataFrame plus the `batchId`, which is the usual way to implement an upsert for update mode. ## Practical guidance Pick append when the consumer wants immutable, final facts and can tolerate watermark latency; update when you want low-latency revisions and the sink can upsert; complete only for a small, dashboard-sized result. Treat the output mode as part of the query definition rather than a runtime knob: for a stateful query it interacts with the checkpointed state layout, so validate before flipping it against an existing checkpoint rather than assuming a restart will absorb the change.

  • Why does a streaming aggregation in append mode require a watermark?
    Append promises a row is final when emitted. A group's aggregate can change whenever a later record joins it, so Spark has no way to know a window is finished unless you declare how late event-time data may be. The watermark supplies that bound, and Spark buffers the window's row in state until the watermark passes the window end. Without it the query is rejected at start.
  • Why does the file sink support only append mode?
    The file sink commits output files and records them in its `_spark_metadata` log; once a file is committed it is not rewritten. Update mode would require replacing rows already inside a committed file, and complete mode would require replacing the entire output every trigger. Neither maps onto immutable file commits, so Spark restricts the sink to append.
  • In complete mode, what happens to aggregation state as the query runs?
    It grows forever. Complete mode must re-emit every group on every trigger, so Spark retains all aggregation state and a watermark does not evict it. That is fine for a bounded key space such as counts by country, and a memory problem for anything user- or session-keyed. If you need eviction, use append or update with a watermark instead.

saying these in an interview costs you the question

  • Says append mode works for a streaming aggregation with no watermark
  • Thinks complete mode writes only the rows that changed
  • Assumes any sink can absorb update-mode output without upserts
  • Believes a watermark evicts state in complete mode
  • Confuses the output mode with the trigger interval

context

open as a page

In Spark Structured Streaming, what does withWatermark do to state and to late rows?

level: middleimportance: must knowfreq 68%

basics

~20 s

withWatermark declares how late event-time data may arrive. Spark tracks the maximum event time it has seen, subtracts that threshold, and uses the result to finalize windows, evict their state, and drop rows older than it.

open as a page

How does a Spark Structured Streaming query achieve end-to-end exactly-once output?

level: seniorimportance: should knowfreq 60%

basics

~20 s

Spark writes each micro-batch's input offset range to a write-ahead log under checkpointLocation before processing it, so a restart replays exactly that batch. Exactly-once then holds if the source is replayable by offset and the sink is idempotent or transactional.

open as a page

A Spark Structured Streaming windowed count in append mode emits nothing for an hour — why?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Append mode holds each window until the watermark passes its end, and the watermark only advances from event times in data actually read. A stalled watermark, an oversized delay threshold, or a misbound withWatermark all produce a healthy query that emits nothing.

open as a page

When would you schedule a Spark Structured Streaming query with Trigger.AvailableNow instead of running it continuously?

level: principalimportance: should knowfreq 36%

basics

~20 s

Use Trigger.AvailableNow when the latency budget is minutes to hours and cluster cost dominates: it processes everything currently available in several rate-limited micro-batches, then stops, keeping checkpointed offsets and exactly-once semantics while the cluster runs only briefly.

open as a page

Why switch a stateful Spark Structured Streaming query to RocksDBStateStoreProvider?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

The default state store keeps every key in the executor's JVM heap, so large state means heap pressure and long GC pauses. RocksDBStateStoreProvider moves state into native memory and local disk, letting state exceed the heap at the cost of serialization and disk I/O.

open as a page