# How Pathway Handles Late-Arriving Data and Out-of-Order Events in Streaming Pipelines

> Pathway tackles late-arriving data and out-of-order events in streaming pipelines with ignore_late and postpone operators

- Repository: [Pathway/pathway](https://github.com/pathwaycom/pathway)
- Tags: how-to-guide
- Published: 2026-03-06

---

**Pathway handles late-arriving data and out-of-order events by routing late records to a dedicated stream via the `ignore_late` operator while buffering events with `postpone` until logical timestamps advance, ensuring monotonicity in Differential Dataflow pipelines.**

Pathway's streaming engine, built on Differential Dataflow, treats timestamps as the primary ordering mechanism to process late-arriving data and out-of-order events without dropping records. When events arrive behind the current watermark for a time column, the system classifies them as *late* and processes them through specialized operators that maintain pipeline correctness. This article examines the specific implementation details in the `pathwaycom/pathway` repository, focusing on the time-column operators in [`src/engine/dataflow/operators/time_column.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/operators/time_column.rs) that manage these edge cases.

## The `ignore_late` Operator: Splitting On-Time and Late Records

The core mechanism for handling late-arriving data resides in the **`ignore_late`** operator, implemented in [`src/engine/dataflow/operators/time_column.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/operators/time_column.rs) (lines 660-678). This operator maintains per-instance state tracking the greatest column-time already emitted (`max_column_times`), enabling it to classify incoming records based on their temporal relationship to the pipeline's current progress.

For each incoming record, the operator extracts two critical values:
- **threshold_time_extractor**: The point after which data is considered on-time
- **current_time_extractor**: Used to update the threshold after processing a batch

If a record's timestamp is **greater than or equal to** the stored maximum for its instance, the operator emits it on the **regular** output stream. Otherwise, it routes the record to a **late** output stream. This dual-stream approach keeps late records available for downstream operators that may need to handle them specially, such as replay mechanisms, alerting systems, or selective discarding.

```rust
// src/engine/dataflow/operators/time_column.rs
// Lines 660-678 – implementation of ignore_late
pub fn ignore_late<…>(…)

```

## Buffering Out-of-Order Events with `postpone`

While `ignore_late` separates late records, the **`postpone`** trait (also in [`time_column.rs`](https://github.com/pathwaycom/pathway/blob/main/time_column.rs), lines 651-660) provides the buffering infrastructure that holds records until the appropriate logical time. This operator wraps a collection, extracts a threshold column time, and buffers records until the operator's internal release-threshold advances.

When the threshold time is reached, `postpone` releases buffered records in order, guaranteeing that later-arriving events are only emitted after the appropriate logical time. This mechanism is essential for stateful operators that require a monotonic view of time, including the **freeze** operator which relies on `postpone` to maintain consistent snapshots.

```rust
// src/engine/dataflow/operators/time_column.rs
// Lines 651-660 – call to postpone (used by the "freeze" operator)
ignore_late(
    self,
    threshold_time_extractor,
    current_time_extractor,
    instance_extractor,
)

```

## State Retention and the `forget` Operator

When pipelines need to **expire** old state, Pathway uses the **`forget`** operator, which internally calls `postpone` on the *negative* (retraction) stream. This produces "forgetting" records that are emitted only **after** the current time has passed the retention threshold, ensuring that out-of-order retractions are processed correctly without corrupting downstream state.

The implementation (lines 664-691) demonstrates how Pathway handles the complexity of late-arriving deletions:

```rust
// src/engine/dataflow/operators/time_column.rs
// Lines 664-691 – forget implementation
let forgetting_stream = self.negate().postpone(
    self.scope(),
    … // same extractors as postpone
)?;

```

## Validation Through Integration Testing

The integration test suite in [`tests/integration/test_time_column.rs`](https://github.com/pathwaycom/pathway/blob/main/tests/integration/test_time_column.rs) provides explicit verification of the late-forwarding logic:

- **`test_core_late_forwarding`** (approximately line 190): Verifies that records whose timestamps are behind the current column time are correctly routed to the late stream
- **`test_core_late_forwarding_ignore_retraction`**: Ensures that late **retractions** are also handled without corrupting the output

These tests exercise the `ignore_late` path and confirm that late data does not break the monotonicity guarantees of the pipeline.

```rust
// tests/integration/test_time_column.rs
// Example test reference
fn test_core_late_forwarding() { … }

```

## Practical Implementation Example

To implement late-data handling in your own Pathway pipelines, you combine the `postpone` operator for buffering with `ignore_late` for stream separation:

```rust
use pathway::prelude::*;

// Example: a collection with a time column `event_ts`
let source = stream
    .map(|event| (event.id, event))
    .postpone(
        stream.scope(),
        |(_, e): &(_, Event)| e.event_ts,      // threshold extractor
        |(_, e): &(_, Event)| e.event_ts,      // current time extractor
        |(_, e): &(_, Event)| e.instance_id,   // instance extractor
        flush_on_end: true,
        update_time_before_emitting: false,
        |coll| Ok(coll)                         // no extra logic
    )?;

// Split late records
let (on_time, late) = ignore_late(
    source,
    |(_, e): &(_, Event)| e.event_ts,          // threshold
    |(_, e): &(_, Event)| e.event_ts,          // current
    |(_, e): &(_, Event)| e.instance_id,       // instance
);

```

The first snippet buffers records until the threshold time is reached, while the second demonstrates extracting late records into a separate stream for specialized processing.

## Summary

- **Late detection**: The `ignore_late` operator maintains per-instance maximum column times and splits late records into a dedicated stream separate from on-time data.
- **Event buffering**: The `postpone` trait (utilized by `ignore_late`, `freeze`, and `forget`) buffers records until configurable threshold times are reached, guaranteeing correct ordering for stateful operations.
- **Retention management**: The `forget` operator builds on `postpone` to emit "forget" records only after retention cutoffs, safely handling out-of-order deletions without state corruption.
- **Verified correctness**: Integration tests in [`tests/integration/test_time_column.rs`](https://github.com/pathwaycom/pathway/blob/main/tests/integration/test_time_column.rs) validate that late data and retractions maintain pipeline monotonicity.

## Frequently Asked Questions

### What happens to late-arriving records in Pathway?

Pathway does not drop late-arriving records. Instead, the `ignore_late` operator compares each record's timestamp against the per-instance `max_column_times` and routes late events to a separate output stream. This allows downstream operators to choose whether to replay, alert on, or discard late data while the main pipeline continues processing on-time records.

### How does Pathway ensure monotonicity with out-of-order events?

The **`postpone`** operator ensures monotonicity by buffering records until the logical time (extracted via `threshold_time_extractor`) reaches a release threshold. Records are emitted in timestamp order only when the pipeline's current time advances past their event timestamps, preventing stateful operators from seeing time move backwards.

### Can late data be recovered or processed differently?

Yes. Because `ignore_late` splits streams into `(on_time, late)` tuples, you can attach distinct processing logic to the late stream. This enables use cases such as writing late records to a side table for audit purposes, triggering alerts, or attempting to merge late updates into already-emitted results.

### How does Pathway handle late retractions?

The **`forget`** operator processes retractions by calling `postpone` on the negative (retraction) stream. This ensures that "forget" records are emitted only after the retention threshold has passed, preventing out-of-order deletions from incorrectly removing state that was updated by later-arriving on-time records.