How to Implement Stateful Transformations Using Pathway Operators: Joins and Aggregations
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, 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), 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 (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.
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, 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.
# 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. Implement the api.CombineMany interface to define how values combine and how state persists across batches.
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:
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
JoinContextandGroupedContextobjects 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.pyfor join logic,pathway/internals/groupbys.pyfor grouping mechanics, andpathway/internals/reducers.pyfor aggregation definitions. - Custom state is supported via
StatefulManyReducerinpathway/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. 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) 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. 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.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →