Event-time processing groups records by when they occurred, while a watermark states how much out-of-order arrival a live aggregate is prepared to wait for.
Event time and late arrivals: close windows with an explicit correction policy
Two clocks answer different questions
A receipt submitted at 23:58 may reach the pipeline after midnight. Event-time reporting assigns it to the submission day; ingestion-time reporting assigns it to the day it arrived. Neither is automatically correct. State the business question, record both clocks, and measure the delay distribution. A watermark such as maximum observed event time minus a configured lateness allowance tells the processor which windows it considers complete.
Do not confuse complete with final
A watermark is a processing policy, not evidence that no older event will ever arrive. Events later than the allowed delay need an explicit route: revise prior aggregates, publish a restatement, or quarantine them with a manual reconciliation queue. Dropping them without a count changes historical rates silently. Cohort metrics] need the same late-arrival rule in reports.
Bound state and correction cost
Holding open windows consumes memory proportional to active keys and windows. A long allowed lateness period improves completeness but delays publication and enlarges state. A short period improves freshness but causes more corrections. Measure late-event rate, aggregate revisions and state size rather than picking an arbitrary delay.
Test across the boundary
Feed events in this order: a timely receipt, a newer-day receipt, then a late receipt for the prior day. Verify the prior-day total and correction signal under the chosen policy. Also test duplicate delivery and timezone conversion. A window keyed by local midnight can shift during daylight-saving transitions.
Implementation
from datetime import timedelta
latest_event_time = max(event.event_time for event in received_events)
watermark = latest_event_time - timedelta(hours=3)
for event in received_events:
if event.event_time < watermark:
late_events.append(event)
else:
open_windows[event.event_time.date()].append(event)Performance and operating cost
Each event can be assigned in O(1) expected time, but open state grows with active windows and keys. Corrections may require rebuilding downstream aggregates even when ingestion itself is cheap.
Common Mistakes
- Do not use ingest time when the metric is defined by event time.
- Do not drop records past a watermark without an explicit count and repair path.
- Do not describe an event-time window as immutable before correction policy closes it.
Read next
- Backfills: rebuild history without exposing a half-written result
- Data quality gates: quarantine bad rows and reconcile complete batches
- Metric denominators and cohorts: make a rate reproducible
- Data source contracts: preserve raw records before transformation
Continue the workflow: Training-serving parity: compare feature values at one prediction clock.
Continue the workflow: Event-time windows: place each event by occurrence, not arrival.
