How to Handle Late-Arriving Data and Watermarking in Streaming Systems: A Flink Implementation Guide

Handle late-arriving data in streaming systems by defining watermarks that lag behind event time by a fixed interval, allowing Apache Flink to trigger window computations only after the grace period expires.

Streaming pipelines frequently encounter events that arrive out of order or behind schedule, complicating accurate windowed aggregations. The Data Engineer Handbook repository demonstrates a production-ready pattern for late-arriving data and watermarking in streaming systems using Apache Flink's event-time processing capabilities. By implementing watermarks in aggregation_job.py, you can maintain result accuracy while tolerating bounded lateness in your Kafka streams.

Understanding Event Time vs. Processing Time

In distributed streaming architectures, event time (when an event actually occurred) differs fundamentally from processing time (when the system receives it). Without watermarks, windows would close based on processing time, risking dropped or incorrectly assigned late records. A watermark is a timestamp that signals the system's confidence that all events up to that time have been observed.

According to the DataExpert-io/data-engineer-handbook source code, Flink uses watermarks to determine when a window is complete. When the watermark passes the window's end timestamp, Flink triggers the computation and emits results, ensuring that late-arriving data within the grace period is included.

The Flink job in intermediate-bootcamp/materials/4-apache-flink-training/src/job/aggregation_job.py implements a three-step pattern for handling late data: extracting event time, defining the watermark strategy, and applying it to windowed aggregations.

Extract Event-Time from Raw Payload

First, convert the raw string timestamp to a proper TIMESTAMP type. This creates the event-time column that serves as the foundation for watermarking.

CREATE TABLE process_events_kafka (
    ip VARCHAR,
    event_time VARCHAR,
    host VARCHAR,
    url VARCHAR,
    -- Convert ISO-8601 string to timestamp for event-time processing
    window_timestamp AS TO_TIMESTAMP(event_time, 'yyyy-MM-dd''T''HH:mm:ss.SSS''Z'''),
    WATERMARK FOR window_timestamp AS window_timestamp - INTERVAL '15' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'events_topic',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json'
);

The TO_TIMESTAMP function parses the ISO-8601 string, while the WATERMARK clause defines a 15-second grace period. This tells Flink to expect events up to 15 seconds older than the maximum observed timestamp before considering a window complete.

Configure the Grace Period

The interval INTERVAL '15' SECOND represents the maximum out-of-orderness the system tolerates. You can tune this value based on your source characteristics—higher values increase latency but improve completeness, while lower values reduce latency but risk dropping late events.

Windowing with Watermarked Streams

Once watermarks are defined, apply them to tumbling windows that aggregate data into fixed-time buckets. The watermark ensures the window waits for late data before triggering.

Tumbling Window Aggregation

The following Python code from aggregation_job.py groups events into 5-minute windows using the watermarked column:

from pyflink.table import EnvironmentSettings, TableEnvironment
from pyflink.table.expressions import col, lit
from pyflink.table.window import Tumble

# Assuming t_env is the TableEnvironment

t_env.from_path("process_events_kafka") \
    .window(Tumble.over(lit(5).minutes).on(col("window_timestamp")).alias("w")) \
    .group_by(col("w"), col("host")) \
    .select(
        col("w").start.alias("event_hour"),
        col("host"),
        col("host").count.alias("num_hits")
    ) \
    .execute_insert("aggregated_results_table")

The Tumble.over(lit(5).minutes).on(col("window_timestamp")) declaration instructs Flink to wait until the watermark exceeds each 5-minute bucket's end time before computing the COUNT aggregation.

Handling Late Data Beyond the Watermark

Events arriving after the watermark (more than 15 seconds late) are considered late data. By default, Flink drops these records, but the system supports alternative strategies:

  • Side outputs: Route late events to a separate stream for manual inspection or reprocessing using getSideOutput().
  • Allowed lateness: Extend the window state to update results if late data arrives within an additional grace period (configured via allowedLateness).

The repository's start_job.py—located at intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py—provides a useful contrast by reading from Kafka without watermarking, demonstrating why explicit event-time handling is necessary for accurate results.

Summary

  • Event-time semantics align processing with actual occurrence timestamps rather than ingestion time, ensuring accuracy despite network delays.
  • Watermarks provide deterministic triggers for window completion, balancing latency against data completeness.
  • The 15-second lag in aggregation_job.py demonstrates a configurable grace period for bounded out-of-order data.
  • Tumbling windows use watermarked columns to ensure aggregations include all expected events before emitting results.
  • Late data strategies include side outputs and allowed lateness configurations for events exceeding the watermark threshold.

Frequently Asked Questions

What is the difference between event time and processing time?

Event time refers to the timestamp embedded in the data itself (when the event actually occurred), while processing time is the wall-clock time of the machine executing the stream processor. Event time ensures accurate results despite network delays or out-of-order delivery, whereas processing time offers lower latency but potential inaccuracy when handling late-arriving data.

How do I choose the right watermark lag interval?

Select a lag interval that accommodates your source's maximum expected delay. Analyze historical data to find the 99th percentile of arrival delays, then add buffer time. The 15-second example in the Data Engineer Handbook suits well-ordered sources; chaotic sources with unpredictable network conditions may require minutes or even hours of lag.

What happens to data that arrives after the watermark?

Events arriving after the watermark has passed are classified as late data. By default, Flink discards these records from the main computation. However, you can configure side outputs to capture them separately for audit trails or use allowed lateness to update previously emitted window results within a secondary grace period.

Can watermarks be used with sliding windows?

Yes, watermarks work with any window type including sliding windows, session windows, and tumbling windows. The watermark triggers computation whenever it advances past a window's end timestamp, regardless of the windowing strategy. Sliding windows simply create overlapping buckets that may trigger independently as the watermark progresses.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →