skip to content

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

level: middleimportance: must knowfreq 68%

answer

  1. a promise about how late is late enough
  2. derived from the data, never the clock
  3. it moves only forward
  4. it both closes windows and frees memory
  5. max event time seen minus the threshold

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.

solid answer

~50 s

`withWatermark("eventTime", "10 minutes")` tells Spark that records may lag the newest event time by up to ten minutes. Spark computes the watermark at the end of each micro-batch as *(maximum event time observed so far) − (threshold)*, and applies that value in the **next** batch, so it always lags one trigger behind. Stateful operators use it for two things: they finalize and emit a window once the watermark passes the window's end, and they then evict that window's state, which is what keeps a streaming aggregation's memory bounded. Records whose event time falls below the current watermark are dropped by those operators. The guarantee is one-directional: anything delayed *less* than the threshold is processed, anything delayed *more* may or may not be. With several input streams Spark takes the minimum of their watermarks by default.

code

python · 10 lines
python
counts = (events
    .withWatermark("eventTime", "10 minutes")   # before the aggregation
    .groupBy(window("eventTime", "5 minutes"), "deviceId")
    .count())

# WRONG: watermark declared after the aggregation has no effect
bad = (events
    .groupBy(window("eventTime", "5 minutes"), "deviceId")
    .count()
    .withWatermark("eventTime", "10 minutes"))

go deeper

for a junior

Know the call shape — withWatermark on the event-time column with a delay like "10 minutes" — and that it exists so Spark can decide a window is finished and stop holding its data in memory.

for a middle

Explain the arithmetic: maximum observed event time minus the threshold, recomputed each micro-batch and applied in the next. Be able to say what it does for windows, for state eviction, and for late rows.

for a senior

Show you size the threshold from measured arrival lag and that you monitor it. Expect to be asked why output stopped and to reach for the watermark value and numRowsDroppedByWatermark in lastProgress before touching any config.

for a principal

Own the completeness-versus-latency contract this sets for the whole pipeline, including what happens to records you knowingly drop — whether a late-arriving-facts reconciliation path exists downstream, and how the multi-source minimum policy couples independent teams' latency.

## Why a watermark exists at all A streaming aggregation has to answer a question a batch job never faces: *when is a group finished?* A batch `GROUP BY` reads the whole input and then emits, but a stream has no end. If Spark kept every window's partial aggregate forever, state would grow without bound; if it closed windows eagerly, records arriving a little late would be lost or would produce a second, contradictory answer. The watermark is the explicit lateness contract that resolves this. `withWatermark(colName, delayThreshold)` says: *for this event-time column, I am willing to wait `delayThreshold` past the newest event time I have seen; beyond that, I accept losing records.* Everything else — when windows emit, when state is released, which rows are dropped — follows from that one declaration. ## How the value is computed The watermark is derived from the data, not from the wall clock. At the end of a micro-batch Spark takes the maximum value of the watermarked event-time column across every row it has processed so far and subtracts the threshold. That value is then used by stateful operators during the *following* micro-batch, so the watermark always trails the data by one trigger. Two consequences matter in practice. First, the watermark only ever moves forward — it is monotonic, so a burst of very old records cannot drag it backwards. Second, it advances only from data the query actually reads: if the stream goes quiet, the maximum event time stops moving, the watermark freezes, and buffered windows never close. A stalled watermark on an idle source is one of the most common "my streaming job stopped producing output" reports. ## Where to call it `withWatermark` must be applied to the event-time column *before* the stateful operator that should use it, and the same column has to be the one the aggregation groups on (directly, or via `window()`), otherwise the watermark has nothing to bind to. Applying it after the aggregation, or on a different column than the one being windowed, produces a query where the watermark simply has no effect — the query runs, state grows, and nothing is ever dropped or evicted. ## What it does for each operator **Windowed aggregations.** A window is finalized when the watermark passes the window's end. In append mode the window's row is emitted at that moment and its state is released; in update mode revisions are emitted continuously and the state is released at the same point. Effective output latency is therefore roughly *window length + delay threshold*. **Arbitrary stateful processing.** `flatMapGroupsWithState` with `GroupStateTimeout.EventTimeTimeout()` fires a group's timeout callback when the watermark passes the timestamp you registered, which is how you expire sessions. **Stream–stream joins.** Watermarks on both sides plus a time-range condition on the join let Spark bound how long unmatched rows must be buffered. Outer joins require this, because the NULL-padded row can only be emitted once the watermark proves no match will ever arrive. **Deduplication.** `dropDuplicates` on a watermarked stream keeps the seen-keys set bounded by the watermark instead of forever. ## Multiple input streams When a query reads more than one stream, each source contributes its own watermark and Spark reduces them to a single global value. The default policy is the **minimum**, controlled by `spark.sql.streaming.multipleWatermarkPolicy` (`min` by default, `max` available). Taking the minimum is the safe choice — it means a slow or backfilling source holds the whole query's watermark back rather than causing its counterpart's records to be dropped. It also means one lagging input can stall output for everything downstream, which is worth knowing before you union a live topic with a historical replay. ## The guarantee is asymmetric Spark is careful about what it promises. A record delayed by less than the threshold **is** processed. A record delayed by more than the threshold **may or may not** be dropped — the more delayed it is, the more likely it is discarded, but there is no promise that a slightly-over-threshold record is thrown away. Do not build correctness on "anything late is definitely dropped"; build on "anything within the threshold is definitely counted." ## Observing it `query.lastProgress` exposes an `eventTime` block containing the current `watermark`, and each entry in `stateOperators` reports `numRowsTotal`, `numRowsUpdated` and `numRowsDroppedByWatermark`. Those three numbers answer most watermark questions in production: whether the watermark is advancing at all, how big state has grown, and how much data you are actually losing to lateness. A steadily climbing `numRowsTotal` with a frozen `watermark` is the signature of a watermark that is not doing its job. ## Choosing the threshold The threshold is a latency-versus-completeness dial. Larger means fewer dropped records, more retained state, and later output in append mode; smaller means faster, cheaper, lossier. Size it from a measured distribution of *(arrival time − event time)* on the real source rather than from intuition — a 99th-percentile lag with a margin is the usual starting point.

  • If the source goes idle, what happens to the watermark and to buffered windows?
    The watermark advances only from the maximum event time in data the query reads, so an idle source freezes it. Windows already buffered in state never pass the watermark and never emit, and the query looks stalled even though it is healthy. Confirm with the `watermark` field in `lastProgress`; the fix is upstream traffic, a smaller threshold, or a design that does not depend on a quiet stream closing windows.
  • What does Spark actually guarantee about records delayed by more than the threshold?
    Nothing definite. The guarantee runs one way: records delayed less than the threshold are processed. Records delayed more may or may not be dropped, and the further past the threshold they are, the more likely they are discarded. Build correctness on the inclusion guarantee, never on the assumption that late data is reliably excluded.
  • With two watermarked input streams, which watermark do stateful operators use?
    By default the minimum of the two, controlled by `spark.sql.streaming.multipleWatermarkPolicy`. The minimum is conservative: a lagging source holds the global watermark back so its counterpart's records are not dropped prematurely. The trade-off is that one slow input stalls output for the whole query, which matters when you union a live stream with a historical backfill.

saying these in an interview costs you the question

  • Says the watermark is based on the wall clock or processing time
  • Claims every record older than the watermark is guaranteed to be dropped
  • Calls withWatermark after the aggregation and expects it to apply
  • Thinks the watermark can move backwards when old data arrives
  • Assumes an idle source keeps the watermark advancing

context