Watermarks and late data
Event time vs processing time, how watermarks bound state, and what happens to late rows.
On this page
You will learn
- The difference between event time and processing time
- How a watermark is computed and what it promises
- What happens to rows that arrive after the watermark
- How watermarks bound state for windows, joins and dropDuplicates
withWatermark("event_time", "10 minutes") means "once I have seen an event at 12:30, I will not wait for events before 12:20". Spark uses it to finalise windows and to delete old state. Rows older than the watermark may be dropped, so the delay is a trade-off between completeness and latency and memory.Event time vs processing time
Event time is when something happened (the timestamp in the event). Processing time is when Spark sees it. A phone that was offline sends its 12:05 click at 12:40. Counting clicks per 5 minutes of processing time would put it in the wrong bucket; event-time windows put it where it belongs, but only if Spark still has that window open.
counts = (clicks .withWatermark("event_time", "10 minutes") .groupBy(F.window("event_time", "5 minutes"), "page") .count()) counts.writeStream.outputMode("append").format("delta") \ .option("checkpointLocation", chk).toTable("gold.page_clicks")
How the watermark moves
After each micro-batch, Spark computes watermark = max(event_time seen so far) - delay. It only moves forward. A window is final once the watermark passes its end.
Batch 1 · events 12:01, 12:07
Batch 2 · events 12:16, 12:03
Batch 3 · the window closes
Batch 4 · event 12:02 arrives
numRowsDroppedByWatermark in the query progress counts such rows.Why it matters: state
| Operation | State without a watermark | With a watermark |
|---|---|---|
| Windowed aggregation | Every window ever seen, forever | Windows older than the watermark are dropped |
| dropDuplicates | Every key ever seen | Use dropDuplicatesWithinWatermark (Spark 3.5+) or include the event time column |
| Stream-stream join | Both sides buffered forever | Watermarks on both sides plus a time-range join condition let Spark discard old rows |
A stream that runs for months without a watermark grows its state until it runs out of memory or the checkpointcheckpoint: In streaming, the folder where Spark records which input it has processed and its state, so a restarted query continues where it stopped. Learn more → becomes huge. RocksDB state store helps hold large state on disk, but does not remove the need to bound it.
Choosing the delay
- Measure lateness: the distribution of processing time minus event time. Cover the 99th or 99.9th percentile.
- Longer delay: more complete results, but windows finalise later (append mode) and more state is kept.
- Shorter delay: fresher output and less memory, more rows dropped.
- If late data must never be lost, also land raw events in a table and correct aggregates with a periodic batch job.
Watermarks and output modes
In append mode a window is written once, after it is final, so output is delayed by the window length plus the watermark delay. In update mode a window is written every time its count changes, so results are fresh but the sink must handle updates (MERGE). Spark 4.0 adds transformWithState for custom stateful logic with timers and TTL.
Common mistakes
No watermark on a long-running aggregation
Watermark on a different column than the window
Expecting append mode to write windows immediately
Assuming nothing is lost
Key takeaways
- Watermark = max event time seen - delay; it only moves forward.
- Windows older than the watermark are finalised and their state removed.
- Data later than the watermark may be dropped.
- Pick the delay from measured lateness; back it up with batch correction if needed.
Check yourself
3 questions1. Max event time seen is 10:40 and the delay is 15 minutes. What is the watermark?
Show the answer
10:25. 10:40 - 15 minutes.
2. In append mode, when is a window written?
Show the answer
When the watermark passes the window end. Append writes rows that will not change again.
3. What is the main purpose of a watermark besides handling lateness?
Show the answer
Bounding the state Spark must keep. Without it, state for windows, dedup and joins grows forever.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 14 of 24 · 5 min readData skew and salting
Why one task runs for an hour while the rest finish in seconds, and how salting spreads a hot key.
Related lessons
Structured Streaming basicsPySpark functions · 3 min read
dropDuplicatesSpark internals · 4 min read
The memory model
Previous: Structured Streaming basics
Primary sources: Handling late data and watermarking · Window operations on event time