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

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 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 (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.

// 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, 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.

// 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:

// 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 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.

// 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:

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 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.

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 →