A Spark Structured Streaming windowed count in append mode emits nothing for an hour — why?
answer
- healthy query, empty sink
- output waits on something that must move
- that something only moves when data arrives
- two inputs, and the slower one wins
- check lastProgress before touching config
basics
~20 sAppend 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.
solid answer
~50 sIn append mode a windowed aggregate is buffered in state and released only when the watermark moves past the window's end, so the first output is roughly *window length + delay threshold* of **event time** after the window opens — a one-hour window with a 30-minute watermark simply will not emit for 90 minutes. Beyond that expected latency, three things stall it for real. The watermark is computed from the maximum event time in data the query reads, so an idle or very low-volume source freezes it and no window ever closes. With two inputs the default `min` policy means a lagging source holds the global watermark back. And if `withWatermark` was applied to a different column than the one being windowed, or after the aggregation, it binds to nothing. Read `query.lastProgress`: the `eventTime.watermark` value and `stateOperators[].numRowsTotal` tell you which case you have in one look.
code
text · 12 lines{
"batchId" : 412,
"numInputRows" : 0,
"eventTime" : {
"watermark" : "2026-03-01T09:12:00.000Z"
},
"stateOperators" : [ {
"numRowsTotal" : 1843277,
"numRowsUpdated" : 0,
"numRowsDroppedByWatermark" : 0
} ]
}go deeper
Remember that append mode does not emit a window until the watermark has moved past its end, so a streaming aggregation is expected to be quiet at first. Silence for a few minutes is normal, not a failure.
Be able to compute the expected first-output latency as window length plus watermark delay in event time, and to explain that the watermark advances only from event times in data the query actually reads.
Demonstrate the diagnosis order: expected latency first, then lastProgress metrics, then the watermark's binding and the multi-source policy. Interviewers want to see you rule out a healthy-but-slow query before you change any setting.
Own the guardrails that stop this recurring: timestamp sanitisation at ingest so one bad record cannot poison a watermark, a policy against unioning backfills with live feeds, and streaming SLOs expressed in event-time lag rather than job uptime.
## First, is it actually broken? Append mode's contract is that a row is emitted only when it is final. For a windowed aggregation, finality is decided by the watermark: the window's row leaves state when the watermark passes the window's end. So the *designed* latency is approximately `window length + watermark delay threshold` measured in **event time**, not wall clock. A ten-minute tumbling window with a one-hour watermark produces its first row about seventy minutes of event time after the stream starts. A great many "nothing is coming out" reports are this and nothing else, and the fix is a smaller threshold or update mode, not a bug hunt. Establish the expected latency before investigating anything. ## The watermark is frozen If the expected latency has passed and output is still empty, the next question is whether the watermark is moving at all. Spark derives it from the maximum value of the event-time column across rows it has processed, minus the threshold. Nothing else feeds it — not the wall clock, not the trigger, not an empty micro-batch. So a stream that goes quiet freezes the watermark exactly where it was. Batches keep running, the query reports itself healthy, and every window already buffered in state sits there indefinitely. Low-volume topics show the same symptom in a milder form: the watermark advances in jumps whenever a record happens to arrive, and output arrives in bursts rather than steadily. This is a genuine and deliberate difference from engines that emit idleness markers; a Spark watermark simply has no source of progress other than data. ## A second input is holding it back When the query reads more than one stream, Spark reduces the per-source watermarks to one global value, and `spark.sql.streaming.multipleWatermarkPolicy` defaults to `min`. That is the safe policy — the lagging source keeps the watermark low so its own records are not dropped — but it also means a historical backfill unioned with a live topic pins the watermark to the backfill's event times. The live side's windows will not close until the backfill catches up. Setting the policy to `max` unblocks output at the cost of dropping the slow source's data as late; usually the better answer is to not mix a replay and a live feed in one query. ## The watermark is bound to nothing `withWatermark` only takes effect if it is applied to the event-time column *before* the stateful operator, and the aggregation must actually group on that column — directly, or through `window()`. Three misplacements produce a silently inert watermark: - calling `withWatermark` *after* the `groupBy`, so the aggregation was planned without it; - watermarking column `ingestTs` while windowing on `eventTs`; - projecting the event-time column away (or renaming it) between the watermark and the aggregation. In the first case the query would normally have been rejected — append plus a streaming aggregation with no usable watermark fails at analysis — so a running query usually means the watermark exists but is bound to the wrong column, which is worse: state grows and nothing closes. ## Event time is really processing time A subtler variant: the column being watermarked is populated at ingestion, so it always equals roughly "now." That looks fine and the watermark advances normally — but if some upstream stage stamps a constant, or a schema change turned the column into nulls or epoch-zero values, the maximum event time can sit at an absurd value and the watermark with it. A single record carrying a far-future timestamp does the opposite damage: because the watermark is monotonic and derived from the maximum, one bad record from the year 2999 pushes the watermark forward permanently and Spark then drops **all** subsequent real data as late. Symptomatically that shows up as output that stops forever after one burst, with `numRowsDroppedByWatermark` climbing. ## Reading the evidence `query.lastProgress` (and the `StreamingQueryListener` events behind it) answers all of this without guesswork: - `eventTime.max` and `eventTime.watermark` — is the watermark advancing, is it wildly ahead of reality? - `stateOperators[].numRowsTotal` — is state accumulating with nothing ever leaving? - `stateOperators[].numRowsDroppedByWatermark` — is data being discarded as late? - `numInputRows` — is the query reading anything at all? A frozen `watermark` with rising `numRowsTotal` means the source stalled or a second input is holding it back. A far-future `watermark` with rising `numRowsDroppedByWatermark` means a poisoned timestamp. Rising `numRowsTotal` with a healthy, advancing watermark and no output points at the watermark being bound to the wrong column. ## Remedies Shorten the delay threshold if the latency is simply the design. Switch to update mode if the consumer can absorb revisions and you want output before the window closes — it emits partial counts immediately and does not wait on the watermark. Split a backfill out of the live query rather than fighting the `min` policy. Sanitise timestamps at the source so one malformed record cannot advance the watermark past your real data. And in every case, validate against `lastProgress` rather than by staring at an empty output directory.
- How would you confirm the diagnosis from the query's own metrics?Read `query.lastProgress`. A frozen `eventTime.watermark` with rising `stateOperators[].numRowsTotal` means the watermark is not advancing — an idle source or a lagging second input. A far-future watermark with rising `numRowsDroppedByWatermark` means a poisoned timestamp pushed it past your real data. A healthy advancing watermark with growing state and no output points at a watermark bound to the wrong column.
- A single record with a year-2999 timestamp arrives. What happens to the query?The watermark is the maximum observed event time minus the threshold and it never moves backwards, so that one record advances it permanently past every real event time. Every subsequent record is below the watermark and gets dropped as late, and `numRowsDroppedByWatermark` climbs while output stops. Recovery means sanitising or filtering timestamps at ingest and restarting from a checkpoint whose watermark predates the bad record.
- If the consumer cannot wait for the watermark, what changes?Switch to update mode. It emits every group whose value changed in the trigger, so partial counts appear seconds after the window opens instead of waiting for the window to close. The cost is that the sink now receives repeated revisions of the same key and must upsert. Keep the watermark anyway — in update mode it still bounds state even though it no longer gates emission.
saying these in an interview costs you the question
- Assumes append mode emits as soon as the micro-batch finishes
- Expects the watermark to advance on wall-clock time while the source is idle
- Increases the trigger frequency to make windows close sooner
- Thinks an empty micro-batch still pushes the watermark forward
- Never looks at lastProgress before changing configuration