A watermark is a progress estimate for event time; it trades waiting for late data against timely output and bounded state.
Watermarks and late events: define when a window becomes final
Choose a delay budget
If most receipt events arrive within four minutes but mobile uploads can take longer, a four-minute lateness budget produces timely numbers with a known late-event tail. Measure the observed delay distribution before setting it. The watermark does not guarantee no older event will ever arrive. Window boundaries determine which aggregate a late event would change.
Decide the late path
A late event can be dropped with a count, routed to a correction stream, or allowed to revise a previously emitted window while state remains. Each choice changes the consumer contract. A dashboard should label provisional versus final values and version corrections. Atomic publication can prevent readers from mixing old and new window versions.
Handle idle partitions
A slow or idle input partition can hold back a combined watermark even when other partitions are current. Track per-partition progress and define idleness carefully; marking an active but delayed partition idle can make its events unexpectedly late. A single global “lag” number can hide this state.
Rehearse a late record
Emit a 09:58 receipt before the 10:00 hourly window is finalized, then replay the same event after the chosen lateness boundary. The first should update the provisional aggregate; the second should take the declared late path rather than silently changing a final number. Check duplicate IDs in both paths.
Implementation
def lateness_route(event_minute, window_end_minute, watermark_minute, grace_minutes):
if grace_minutes < 0 or window_end_minute <= event_minute:
raise ValueError("invalid window or grace period")
if watermark_minute < window_end_minute:
return "provisional_window"
if watermark_minute < window_end_minute + grace_minutes:
return "revision_window"
return "late_correction_queue"Performance and operating cost
Routing is O(1) per event; holding state for more grace time increases memory roughly with active keys times retained windows. Measure late-event frequency and correction cost rather than picking a watermark delay only from desired dashboard speed.
Common Mistakes
- Do not treat a watermark as proof that older events cannot arrive.
- Do not revise a final metric without a versioned correction contract.
- Do not ignore idle or slow partitions when diagnosing watermark lag.
