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

> Master late-arriving data and watermarking in streaming systems with Apache Flink. Implement effective strategies to ensure accurate event time processing and reliable window computations.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: how-to-guide
- Published: 2026-08-07

---

**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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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.

## Implementing Watermarks in Flink

The Flink job in [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/aggregation_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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.

```sql
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/aggregation_job.py) groups events into 5-minute windows using the watermarked column:

```python
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/start_job.py)—located at [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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.