A streaming engine determines whether an event is late by comparing its event timestamp against the current watermark. When an event arrives three hours late, you can update existing window calculations if it falls within allowed lateness, discard the record completely, or route it to a side output dead-letter table. Choosing between these paths depends on downstream business tolerance for data revisions versus infrastructure processing costs, and teams track the late-event arrival rate as an operational metric.
How watermarks identify late events
The watermark acts as an assertion of completeness, signaling to operators that all events prior to a specific timestamp have arrived. When an incoming record bears a timestamp older than the current watermark, the engine flags it as late.
If the pipeline configures an allowed lateness interval that spans three hours, the engine reopens the target window and incorporates the late arrival:
SELECT window_start, window_end, COUNT(event_id) AS total_events FROM TUMBLE(TABLE stream_events, DESCRIPTOR(event_time), INTERVAL '1' HOUR) GROUP BY window_start, window_end;
Updating windows requires downstream storage systems capable of handling upserts or merges.
Key operational characteristics of updating open windows:
- Storage sinks must support primary key overwrites, such as Apache Iceberg merge operations, Delta Lake tables, or key-value document stores.
- Downstream dashboards consuming these tables will observe mutating historical aggregates, which can confuse end users if reporting figures change hours later.
Handling events beyond allowed lateness
When an arrival exceeds the allowed lateness threshold, the window state has already been purged from the state backend. At this stage, two primary options remain:
- Drop the record: The engine discards the event immediately. This protects memory and avoids re-triggering downstream calculations. It is acceptable for high-volume metric telemetry where minor omissions do not alter aggregate trends.
- Route to a side output: Engines like Apache Flink route late records to a designated side output stream, while Spark pipelines filter them into a dead-letter landing zone on object storage. Data engineers inspect this late table to debug upstream network disruptions.
Reconciliation and monitoring
A common production design drops or routes late events in real time, then schedules a nightly batch reconciliation job. The batch job re-reads raw immutable event logs from object storage and recomputes daily partition summaries, producing clean historical records.
Always emit a custom metric such as stream_late_events_total to track late arrivals against total pipeline volume. If the percentage of late events spikes above acceptable thresholds, investigate client clock skew or upstream transport delays.