A watermark is the engine's progressive claim: "I believe I have seen (almost) all events with event_time less than T." After that, the engine can close windows, emit results, and free state.
Without watermarks, event-time windows never know when late data is "done," so state grows forever.
events by event_time: 10:01 10:02 10:00(late) 10:03 ... watermark advances ----> 9:55 --> 9:57 --> 10:00 --> ... allowed lateness = 5 min window [10:00, 10:05) closes when watermark passes 10:05 + lateness policy
How fresher systems use it
- Spark Structured Streaming:
withWatermark("event_time", "10 minutes") - Flink: assign timestamps + watermark strategy (periodic / punctuated)
What happens to very late data
Often dropped from that window, sent to a side output, or handled with special late-data logic. You trade completeness for bounded memory.
Interview tip: Watermark ≈ "how late can events be before we give up waiting." Tie it to event time, not wall-clock job start.