# How to Implement Stateful Transformations Using Pathway Operators: Joins and Aggregations

> Learn to implement stateful transformations like joins and aggregations in Pathway. Explore how join contexts and grouped bucket reducers efficiently process streaming data with join, groupby, and reduce operators.

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

---

**Pathway implements stateful transformations through mutable *join contexts* and *grouped bucket reducers* that incrementally update results as new streaming data arrives, using the `join`, `groupby`, and `reduce` operators.**

The `pathwaycom/pathway` repository provides a Python streaming framework where tables are treated as unbounded streams. Unlike batch processing systems, Pathway's core operators maintain internal mutable state—referred to as a *universe*—that evolves with each new row, enabling continuous, low-latency computation for joins and aggregations over live data.

## Architectural Overview of Stateful Operators

Pathway's stateful architecture rests on three interconnected components that preserve computation state across incremental updates.

### JoinContext and Universe Management

In [`pathway/internals/joins.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/joins.py), the `JoinResult` class encapsulates a `JoinContext` that stores a **universe**—the set of keys seen so far—and mapping tables for both left and right inputs. When `join()` is invoked, the `_table_join` method constructs this context, which tracks join modes (`INNER`, `LEFT`, etc.) and flags like `exactly_once`. As new rows append to either side, the context updates only the affected key mappings rather than recomputing the entire join.

### GroupedContext and Bucket State

Aggregations rely on `GroupedTable` (defined in [`pathway/internals/groupbys.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/groupbys.py)), which creates a `GroupedContext` containing its own universe. Each distinct grouping key maps to a bucket (via `hash(key) → bucket`) that holds a **partial aggregation**—the intermediate state of reducers like `sum` or `count`. The `TableReduceDesugaring` mechanism (method `_desugaring` in the same file) transforms reducer expressions into column-wise operations evaluated against these buckets.

### Reducer State Persistence

Standard reducers in [`pathway/internals/reducers.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/reducers.py) (e.g., `UnaryReducer`) are stateless descriptions, but at runtime they materialize as concrete engine objects (e.g., `api.Reducer.INT_SUM`) attached to each bucket. For complex logic, the `StatefulManyReducer` class enables custom Python state via the `api.CombineMany` interface, maintaining arbitrary mutable objects across batches inside each group bucket.

## Implementing Stateful Joins with Pathway

The `join` operator creates a `JoinResult` that maintains live references to both input tables. Use it to correlate rows incrementally as data streams in.

```python
import pathway as pw

# Two source tables (could be live streams)

left = pw.debug.table_from_markdown('''
    id | name   | country
    1  | Alice  | US
    2  | Bob    | CA
    3  | Carol  | US
''')

right = pw.debug.table_from_markdown('''
    id | score
    1  | 85
    2  | 92
    4  | 77
''')

# Stateful inner join on the `id` column

joined = left.join(right, left.id == right.id)      # ← creates JoinResult with context

joined = joined.select(
    user_id=left.id,
    name=left.name,
    score=right.score,
)                                                # materializes output

pw.debug.compute_and_print(joined, include_id=False)

```

Output:

```

user_id | name  | score
1       | Alice | 85
2       | Bob   | 92

```

The `joined` object preserves the **join context**; subsequent appends to `left` or `right` trigger incremental updates only for modified keys.

### Join Modes and Left/Outer Variants

Replace `join` with `join_left`, `join_right`, or `join_outer` to control unmatched row retention. These variants share the same implementation in [`joins.py`](https://github.com/pathwaycom/pathway/blob/main/joins.py), differing only in the `mode` argument passed to `_table_join`. For example, `join_left` passes a left-outer mode that retains left-table rows even when the predicate finds no right-table match.

## Performing Stateful Aggregations

After joining, use `groupby` followed by `reduce` to compute rolling statistics. The `groupby` operator creates a `GroupedTable` with a **grouped context** that allocates one bucket per distinct key value.

```python

# Continue from the `joined` table above

agg = (
    joined
    .groupby(joined.country)               # ← creates GroupedTable with buckets

    .reduce(
        count=pw.reducers.count(),
        avg_score=pw.reducers.sum(joined.score) / pw.reducers.count(),
    )
)

pw.debug.compute_and_print(agg, include_id=False)

```

Output:

```

country | count | avg_score
US      | 1     | 85
CA      | 1     | 92

```

The `reduce` call registers reducer objects within each country bucket. The `sum` reducer updates its internal accumulator incrementally as new rows arrive, avoiding full recomputation of the aggregate.

### Custom Stateful Reducers

For aggregations requiring complex Python state (such as maintaining a top-K list), extend `StatefulManyReducer` from [`pathway/internals/custom_reducers.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/custom_reducers.py). Implement the `api.CombineMany` interface to define how values combine and how state persists across batches.

```python
import pathway as pw
from pathway.internals.custom_reducers import StatefulManyReducer
import pathway.internals.api as api

def top_k_reducer(k: int):
    class TopK(api.CombineMany[int]):
        def __init__(self):
            self.heap = []  # internal mutable state

        def combine(self, value):
            import heapq
            heapq.heappush(self.heap, value)
            if len(self.heap) > k:
                heapq.heappop(self.heap)

        def result(self):
            return sorted(self.heap, reverse=True)

    return StatefulManyReducer(combine_many=TopK)

# Usage in a reduction

top5 = top_k_reducer(5)
agg = (
    joined
    .groupby(joined.country)
    .reduce(
        top_scores=pw.reducers.apply_custom(top5, joined.score)
    )
)

```

The `StatefulManyReducer` wrapper ensures that each group bucket instantiates and preserves its own `TopK` object, enabling arbitrary stateful logic within the streaming aggregation.

## Streaming Pipeline Example

Combine joins and aggregations in a continuous pipeline that processes data incrementally:

```python
import pathway as pw
import pandas as pd

# Simulated streaming sources

source_left = pw.Table()
source_right = pw.Table()

# Continuous stateful pipeline

pipeline = (
    source_left
    .join_left(source_right, source_left.id == source_right.id)
    .groupby(source_left.country)
    .reduce(
        total_score=pw.reducers.sum(pw.this.score),
        users=pw.reducers.count(),
    )
)

# Feed data incrementally; only affected keys recomputed

for batch in pd.read_csv('users.csv', chunksize=1000):
    source_left.append(batch[['id', 'name', 'country']])
    source_right.append(batch[['id', 'score']])
    # At any moment `pipeline` reflects the up-to-date aggregated view

```

The pipeline leverages the **join context** and **grouped reducer buckets** to ensure each new batch updates only relevant keys, maintaining low-latency results over unbounded streams.

## Summary

- **Pathway operators are inherently stateful**, utilizing `JoinContext` and `GroupedContext` objects to maintain *universes* of keys and per-key buckets.
- **Incremental computation** occurs automatically: new rows trigger updates only for affected keys rather than full recomputation of joins or aggregations.
- **Core implementation files** include [`pathway/internals/joins.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/joins.py) for join logic, [`pathway/internals/groupbys.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/groupbys.py) for grouping mechanics, and [`pathway/internals/reducers.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/reducers.py) for aggregation definitions.
- **Custom state** is supported via `StatefulManyReducer` in [`pathway/internals/custom_reducers.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/custom_reducers.py), allowing arbitrary Python objects to persist within group buckets across streaming batches.

## Frequently Asked Questions

### How does Pathway maintain state during a join operation?

Pathway maintains state through the `JoinContext` class in [`pathway/internals/joins.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/joins.py). This context stores a *universe* representing the set of keys encountered and internal mapping tables for both left and right inputs. When new rows arrive, the context updates only the mappings for the specific keys affected, preserving previous computation results and enabling incremental join updates.

### What is the difference between stateful and stateless reducers in Pathway?

Stateless reducers (like those defined in [`pathway/internals/reducers.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/reducers.py)) are declarative descriptions of aggregation operations (e.g., `sum`, `count`). At runtime, these materialize as engine objects inside group buckets. Stateful reducers extend `StatefulManyReducer` and implement the `api.CombineMany` interface, allowing arbitrary Python mutable state (such as lists or dictionaries) to persist within each bucket across incremental updates, supporting complex custom aggregations like top-K or unique counts.

### Can I implement custom aggregation logic in Pathway?

Yes. Create a class implementing `api.CombineMany` from the Pathway internals, then wrap it with `StatefulManyReducer` from [`pathway/internals/custom_reducers.py`](https://github.com/pathwaycom/pathway/blob/main/pathway/internals/custom_reducers.py). This reducer maintains internal state across batches and integrates with the standard `groupby(...).reduce(...)` API. Each group bucket receives its own instance of your reducer class, ensuring isolated state per grouping key.

### How does Pathway handle late-arriving data in stateful transformations?

Pathway's architecture supports late data through its **universe** and bucket-based state management. The `JoinContext` and `GroupedContext` retain historical keys and partial aggregations indefinitely (or until garbage collection policies apply). When late rows arrive, they update the relevant keys in the universe or buckets, triggering incremental recomputation only for the affected outputs while preserving the integrity of the overall stateful transformation.