# How Pathway's Incremental Computation Model Enables Real-Time Data Processing

> Discover how Pathway's incremental computation model achieves real-time data processing by updating only changed data with its efficient Rust-based Differential Dataflow engine.

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

---

**Pathway's incremental computation model processes only data changes rather than reprocessing entire datasets by maintaining state in memory and applying delta updates through its Rust-based Differential Dataflow engine.**

Pathway's data processing engine, available in the `pathwaycom/pathway` repository, leverages an **incremental computation model** built on Differential Dataflow to deliver low-latency results for both batch and streaming workloads. Unlike traditional systems that recompute complete pipelines when new data arrives, the engine maintains the current state in memory and applies only the incremental changes (deltas) that occurred since the last processing cycle.

## The Rust Engine Architecture

At its core, Pathway runs a **scalable Rust engine** that implements Differential Dataflow principles to minimize computational overhead. This architecture keeps collections in memory as *differential* collections storing both positive and negative updates (adds and retractions), while maintaining *frontiers* that track which timestamps have been processed. By tracking these changes at the data structure level, the engine can propagate only the necessary deltas through the computation graph rather than recalculating entire datasets.

## Five Stages of Incremental Processing

The incremental computation model operates through a specific pipeline that minimizes redundant work. According to the Pathway source code, this process involves five distinct stages:

### 1. Data Ingestion via Python API

New records enter the system through input collections that represent the current content of a table. In [`src/python_api.rs`](https://github.com/pathwaycom/pathway/blob/main/src/python_api.rs), the Python API forwards input data to the Rust engine through the `run_with_new_dataflow_graph` function, which initializes the dataflow graph and begins processing.

### 2. Differential State Tracking

Each collection is stored as a *differential* collection capable of tracking both positive and negative updates. The [`src/engine/dataflow/variable.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/variable.rs) file implements this using `differential_dataflow::operators::iterate::Variable` to hold mutable state. The engine maintains *frontiers* that indicate precisely which timestamps have been fully processed, enabling accurate change tracking across the pipeline.

### 3. Delta-Based Operator Execution

Transformations including **filter**, **map**, **join**, and **reduce** are expressed as Dataflow operators that receive only *incremental* updates. The `src/engine/dataflow/operators/` directory contains implementations such as [`stateful_reduce.rs`](https://github.com/pathwaycom/pathway/blob/main/stateful_reduce.rs), [`time_column.rs`](https://github.com/pathwaycom/pathway/blob/main/time_column.rs), and [`prev_next.rs`](https://github.com/pathwaycom/pathway/blob/main/prev_next.rs), which demonstrate how operators recompute only the graph partitions that depend on changed records.

### 4. Work Proportional to Change Size

Because each operator processes a tiny delta instead of the whole dataset, the computational effort scales with the size of the change, not the total data size. As documented in the `pathwaycom/pathway` README, the engine "performs incremental computation" and "only processes the data changes rather than reprocessing the entire dataset," yielding low-latency results even for large historical tables.

### 5. Automatic State Persistence

The engine automatically maintains necessary state—including intermediate aggregates and join indexes—in memory. When requested, it persists this state via the `src/persistence/*` modules, enabling crash recovery without reprocessing past data.

## Real-Time Aggregation Example

The following Python implementation demonstrates how Pathway's incremental computation model handles real-time aggregation without requiring code changes between batch and streaming modes:

```python
import pathway as pw

# 1️⃣ Define a simple schema

class Order(pw.Schema):
    id: int
    amount: float

# 2️⃣ Read a CSV (initial batch) and then keep the table open for streaming

orders = pw.io.csv.read("orders/", schema=Order)

# 3️⃣ Incremental aggregation: sum of amounts per id

total_per_id = orders.reduce(
    total = pw.reducers.sum(orders.amount)
)

# 4️⃣ Export results – the table is kept up‑to‑date automatically

pw.io.jsonlines.write(total_per_id, "totals.jsonl")

# 5️⃣ Run – Pathway will process the existing files **and** any new rows that appear later

pw.run()

```

When new rows are appended to the `orders/` directory, only the affected `id` buckets are recomputed, thanks to Pathway's differential dataflow implementation in [`src/engine/dataflow/variable.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/variable.rs) and the delta-processing logic in [`src/engine/dataflow/operators/stateful_reduce.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/operators/stateful_reduce.rs).

## Performance Benefits of Incremental Computation

Pathway's delta-based architecture delivers specific advantages for modern data pipelines:

- **Real-time updates**: Results refresh immediately as new events arrive, with latency proportional to the size of the change rather than the dataset.
- **Unified batch and streaming**: The same pipeline processes static batch loads and continuous streams without code modifications, leveraging the same incremental operators.
- **Efficient joins and aggregations**: Differential timestamps enable the engine to identify exactly which keys changed, recomputing only those specific partitions rather than entire join indexes.

## Summary

- Pathway's incremental computation model uses a Rust-based Differential Dataflow engine to process only data changes rather than full datasets.
- The system maintains state in memory using `Variable` collections in [`src/engine/dataflow/variable.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/variable.rs) and tracks changes through frontiers.
- Operators in `src/engine/dataflow/operators/` receive only delta updates, recomputing only affected graph partitions.
- The [`src/python_api.rs`](https://github.com/pathwaycom/pathway/blob/main/src/python_api.rs) file bridges Python inputs to the Rust engine via `run_with_new_dataflow_graph`.
- Automatic state management in `src/persistence/*` enables fault tolerance without historical reprocessing.

## Frequently Asked Questions

### How does Pathway's incremental computation model differ from traditional batch processing?

Traditional batch processing systems recompute entire pipelines when new data arrives, leading to latency proportional to total dataset size. Pathway's incremental computation model maintains computation state in memory and applies only delta updates, making processing time proportional to the size of the change rather than the historical data volume.

### What is Differential Dataflow and why does Pathway use it?

Differential Dataflow is a computational model that tracks changes through logical timestamps and maintains state as collections of updates (adds and retractions). Pathway uses this model—implemented in Rust via structures like `differential_dataflow::operators::iterate::Variable` in [`src/engine/dataflow/variable.rs`](https://github.com/pathwaycom/pathway/blob/main/src/engine/dataflow/variable.rs)—to enable efficient incremental updates, allowing operators to process only changed records while maintaining correctness across iterative computations.

### How does Pathway maintain state for fault tolerance?

Pathway automatically keeps necessary intermediate state—including aggregates and join indexes—in memory during computation. The persistence layer in `src/persistence/*` saves this state when configured, enabling the engine to recover from crashes without reprocessing historical data. This state management works transparently with the incremental computation model to ensure exactly-once processing semantics.

### Can the same Pathway pipeline handle both batch and streaming data?

Yes. Pathway's architecture unifies batch and stream processing through the same incremental operators. A pipeline defined using `pw.io.csv.read()` or similar connectors processes existing files as a batch, then continues running to process new files or rows as a stream. The [`src/python_api.rs`](https://github.com/pathwaycom/pathway/blob/main/src/python_api.rs) entry point `run_with_new_dataflow_graph` manages both modes using identical dataflow operators, requiring no code changes to switch between batch historical analysis and real-time streaming.