Pathway Built-in Reduce Operators vs Custom UDFs: Architecture and Performance Differences
Built-in reducers map directly to optimized C++ engine primitives for maximum throughput, while custom UDF reducers wrap Python accumulators through a generic stateful-many wrapper that trades performance for extensibility.
Pathway is a Python stream processing framework that provides two distinct mechanisms for aggregating streaming data: highly optimized built-in reduce operators and extensible custom user-defined functions (UDFs). While both approaches appear similar in the public API, they differ fundamentally in implementation depth, state management, and execution performance. This article examines the architectural distinctions between these aggregation methods based on the pathwaycom/pathway source code.
Implementation Architecture
Built-in Reducers (Engine-Native)
Built-in reducers are pre-implemented classes that map directly to native engine operations. In python/pathway/internals/reducers.py, concrete implementations like SumReducer inherit from UnaryReducer and implement the engine_reducer_unary method to return engine-specific tokens:
class SumReducer(UnaryReducer):
def engine_reducer_unary(self, arg_type: dt.DType) -> api.Reducer:
if arg_type == dt.INT:
return api.Reducer.INT_SUM
elif isinstance(arg_type, dt.Array):
return api.Reducer.array_sum(self.strict)
else:
return api.Reducer.float_sum(self.strict)
The public API in python/pathway/reducers.py provides thin wrappers that instantiate these internal classes:
def sum(arg: expr.ColumnExpression, strict: bool = False) -> expr.ReducerExpression:
return _apply_unary_reducer(SumReducer(name="sum", strict=strict), arg)
Because the engine receives specific reducer tokens like api.Reducer.INT_SUM, the aggregation executes entirely within the Pathway engine without Python-level loops or serialization.
Custom UDF Reducers (Python Wrappers)
Custom reducers are generated at runtime by the udf_reducer decorator defined in python/pathway/internals/custom_reducers.py. This decorator constructs a stateful-many reducer that wraps user-defined accumulator classes:
def udf_reducer(reducer_cls: type[BaseCustomAccumulator]):
def wrapper(*args):
@stateful_many
def stateful_wrapper(packed_state, rows):
# 1. Deserialize previous accumulator state
# 2. Split rows into insertions (count>0) and deletions (count<0)
# 3. Update accumulator via update/retract methods
# 4. Serialize new state for next batch
...
return apply_with_type(extractor, ..., stateful_wrapper(*args))
return wrapper
The wrapper handles deserialization of pickled accumulator state, batch processing through user-implemented methods, and reserialization of the state back to a binary blob for engine storage.
State Management and Performance
Stateless vs Stateful Execution
Built-in reducers maintain minimal state—often primitive counters or registers—managed entirely by the C++ engine core. This allows O(1) per-row overhead with zero Python serialization costs.
Custom UDF reducers are inherently stateful. Users implement a subclass of BaseCustomAccumulator that persists arbitrary Python objects between batches through pickle serialization:
class BaseCustomAccumulator(ABC):
@classmethod
@abstractmethod
def from_row(cls, row: list[api.Value]) -> Self: ...
@abstractmethod
def update(self, other: Self) -> None: ...
@abstractmethod
def compute_result(self) -> api.Value: ...
Each micro-batch triggers Python-level deserialization, method dispatch, and reserialization, creating significant overhead compared to native reducers.
Type Safety and Validation
Built-in reducers enforce compile-time type restrictions through return_type declarations and engine_reducer mappings. For example, sum accepts only numeric or array types, with validation occurring before execution reaches the engine.
Custom UDFs place no automatic type restrictions on the framework side. The accumulator's from_row method determines acceptable inputs, and type validation occurs at runtime based on the implementation logic.
Feature Set and Capabilities
Built-in reducers offer optimized behaviors like skip_nones for null handling and specialized variants such as argmin/argmax with ID column tracking. These leverage internal stateful_single and stateful_many decorators for specific execution modes.
Custom UDF reducers expose unlimited flexibility through optional hooks in BaseCustomAccumulator:
retract: Enables efficient deletion handling without full recomputationneutral: Defines custom empty-group states to avoid materializing first rowssort_by: Guarantees deterministic processing order within batches
If retract is not implemented, the framework falls back to a Counter-based recomputation strategy, which materializes and recalculates entire groups during retractions.
Practical Code Examples
Native Aggregation with Built-in Reducers
import pathway as pw
tbl = pw.debug.table_from_markdown('''
id | value
1 | 10
2 | 20
3 | 30
''')
result = tbl.groupby(tbl.id % 2).reduce(total=pw.reducers.sum(tbl.value))
pw.debug.compute_and_print(result, include_id=False)
This invokes SumReducer, which maps to api.Reducer.float_sum and executes entirely within the engine core.
Custom Weighted Average UDF
import pathway as pw
class WeightedAvgAccumulator(pw.BaseCustomAccumulator):
def __init__(self, weighted_sum, weight):
self.weighted_sum = weighted_sum
self.weight = weight
@classmethod
def from_row(cls, row):
value, w = row
return cls(value * w, w)
def update(self, other):
self.weighted_sum += other.weighted_sum
self.weight += other.weight
def compute_result(self):
return self.weighted_sum / self.weight
weighted_avg = pw.reducers.udf_reducer(WeightedAvgAccumulator)
tbl = pw.debug.table_from_markdown('''
group | value | weight
A | 10 | 2
A | 20 | 1
''')
result = tbl.groupby(tbl.group).reduce(avg=weighted_avg(tbl.value, tbl.weight))
Handling Retractions Efficiently
class SumWithRetract(pw.BaseCustomAccumulator):
def __init__(self, total=0):
self.total = total
@classmethod
def from_row(cls, row):
(value,) = row
return cls(total=value)
def update(self, other):
self.total += other.total
def retract(self, other):
self.total -= other.total
def compute_result(self):
return self.total
sum_reducer = pw.reducers.udf_reducer(SumWithRetract)
Implementing retract avoids the expensive fallback to full group recomputation when streaming deletions occur.
Summary
- Built-in reducers in
python/pathway/internals/reducers.pymap to native C++ primitives viaengine_reducer_unary, offering near-native performance with minimal state overhead. - Custom UDF reducers use the
udf_reducerdecorator inpython/pathway/internals/custom_reducers.pyto wrap Python accumulators, enabling arbitrary aggregation logic at the cost of pickle serialization overhead. - Built-in operators enforce compile-time type constraints and execute statelessly when possible, while custom accumulators manage serialized Python state through
BaseCustomAccumulatormethods. - Custom reducers support advanced streaming features like retractions and deterministic sorting through optional hooks, whereas built-ins provide optimized implementations for common aggregations.
Frequently Asked Questions
When should I use a custom UDF reducer instead of built-in operators?
Use custom UDF reducers when implementing domain-specific aggregations not covered by Pathway's standard library, such as exponential moving averages, weighted statistics, or complex geospatial operations. Built-in reducers are preferred for standard operations like sum, count, min, and max due to their superior performance characteristics.
Why are custom UDF reducers slower than built-in ones?
Custom reducers incur Python serialization overhead because the accumulator state is pickled between batches, and each row batch triggers Python method calls via update or retract. Built-in reducers execute directly in the Pathway engine's C++ core using api.Reducer.* primitives without crossing the Python boundary per row.
Can custom UDFs handle streaming retractions efficiently?
Yes, if you implement the optional retract method in your BaseCustomAccumulator subclass. Without this method, Pathway falls back to a Counter-based recomputation strategy that materializes and recalculates the entire group when deletions occur, which significantly impacts throughput for high-churn streams.
Do built-in reducers support custom neutral elements like UDFs?
No, built-in reducers use implicit neutral elements defined by the engine architecture, such as zero for summation. Custom UDF reducers can define explicit neutral elements via the optional neutral classmethod, allowing you to specify behavior for empty groups without waiting for the first row to arrive in the stream.
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 →