Skip to content
Great engineers know 4 min read · withWatermark 2 practice problems ↓

Watermarks and late data

Event time vs processing time, how watermarks bound state, and what happens to late rows.

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

Read first

Comfortable with these? Read on.

TL;DR A watermark tells Spark how late data may arrive: 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.

PySpark · Clicks per 5-minute window, allowing 10 minutes of lateness
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

Max event time 12:07, so the watermark becomes 11:57. Windows 12:00-12:05 and 12:05-12:10 are open and kept in state. In append mode nothing is written yet.

Batch 2 · events 12:16, 12:03

12:03 is late but newer than the watermark (11:57), so it is counted in 12:00-12:05. Max is now 12:16, so the watermark becomes 12:06.

Batch 3 · the window closes

Because the watermark (12:06) is past 12:05, window 12:00-12:05 is final: append mode writes it once, and its state is deleted.

Batch 4 · event 12:02 arrives

12:02 is older than the watermark. Its window is already final and gone, so the row is dropped. numRowsDroppedByWatermark in the query progress counts such rows.
Note: the guarantee is one-sided. Data later than the watermark is not guaranteed to be counted; it may or may not be dropped depending on timing. Data within the delay is guaranteed to be counted.

Why it matters: state

OperationState without a watermarkWith a watermark
Windowed aggregationEvery window ever seen, foreverWindows older than the watermark are dropped
dropDuplicatesEvery key ever seenUse dropDuplicatesWithinWatermark (Spark 3.5+) or include the event time column
Stream-stream joinBoth sides buffered foreverWatermarks 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

State grows forever.

Watermark on a different column than the window

The watermark must be on the event-time column used for the window or join.

Expecting append mode to write windows immediately

Append writes only final windows, after the watermark passes.

Assuming nothing is lost

Rows later than the delay can be dropped; monitor numRowsDroppedByWatermark.

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 questions

1. 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.

Solve: Late Data in the Raw Zone →

Keep going

Up next · lesson 14 of 24 · 5 min read
Data 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

Previous: Structured Streaming basics

Primary sources: Handling late data and watermarking · Window operations on event time