Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Watermarks and late events: define when a window becomes final

Last updated: 5 Oct 20265 min read
tutorial
IntermediateBy AITrove Editorial

A watermark is a progress estimate for event time; it trades waiting for late data against timely output and bounded state.

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

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

Read next

ai-data
streaming-analytics
Storage details