How Pathway Implements Windowing Operations: Tumbling, Sliding, and Session Windows
Pathway implements tumbling, sliding, and session windows through the unified Table.windowby() API, using specialized classes _SessionWindow and _SlidingWindow to assign rows to temporal buckets before flattening and aggregation.
Pathway is an open-source data processing framework that provides sophisticated temporal analytics capabilities. Its windowing operations—implemented primarily in python/pathway/stdlib/temporal/_window.py—enable stream processing with time-based grouping. This article examines the source code in pathwaycom/pathway to explain how these windowing primitives handle event assignment, merging, and late-arriving data.
The Unified Windowing Architecture
All three windowing primitives share a common six-stage workflow implemented in the windowby() method chain:
-
Construction – Users invoke
pw.temporal.session(),pw.temporal.sliding(), orpw.temporal.tumbling()to instantiate concreteWindowsubclasses with validated parameters. -
Window Assignment – The concrete window’s
_apply()method executes. For sliding and tumbling windows,_window_assignment_function()generates candidate windows and filters valid intervals. For session windows,_merge()iteratively groups adjacent rows based on predicates or maximum gaps. -
Table Enrichment – Hidden columns
_pw_windowand_pw_keyannotate each row with window metadata and original timestamps. -
Flattening – The
flatten()utility expands window lists into separate rows per window while preserving payload data. -
Behavior Handling – Optional temporal behaviors (
delay,cutoff,keep_results) fromtemporal_behavior.pybuffer or freeze data to handle late arrivals. -
Grouping – The system groups by canonical window columns (
_pw_window,_pw_window_start,_pw_window_end,_pw_instance) before applying user-defined reducers.
Session Windows: Dynamic Gap-Based Grouping
Session windows group consecutive events that satisfy a temporal proximity predicate or maximum gap threshold. The _SessionWindow class (lines 64‑78 in _window.py) encapsulates this logic.
Merging Logic and Group Representation
The _merge() method (lines 70‑76) determines session membership by either invoking a user-supplied predicate via pw.apply_with_type or comparing numeric gaps against self.max_gap. When next - cur < self.max_gap evaluates true, rows merge into the same session.
The _compute_group_repr() method (lines 83‑103) creates a pointer column _pw_window that uniquely identifies each session group. This pointer serves as the grouping key during the aggregation phase.
Application Flow
The _apply() method (lines 106‑143) validates arguments, computes group representations, and aggregates per session before returning a GroupedTable. This implementation supports variable-length windows that adapt to event density rather than fixed clock time.
import pathway as pw
t = pw.debug.table_from_markdown(
'''
| instance | t | v
1 | 0 | 1 | 10
2 | 0 | 2 | 1
3 | 0 | 4 | 3
4 | 0 | 8 | 2
5 | 0 | 9 | 4
6 | 0 | 10| 8
7 | 1 | 1 | 9
8 | 1 | 2 | 16
'''
)
result = t.windowby(
t.t,
window=pw.temporal.session(predicate=lambda a, b: abs(a - b) <= 1),
instance=t.instance,
).reduce(
pw.this._pw_instance,
pw.this._pw_window_start,
pw.this._pw_window_end,
min_t=pw.reducers.min(pw.this.t),
max_v=pw.reducers.max(pw.this.v),
count=pw.reducers.count(),
)
pw.debug.compute_and_print(result, include_id=False)
Output:
_pw_instance | _pw_window_start | _pw_window_end | min_t | max_v | count
0 | 1 | 2 | 1 | 10 | 2
0 | 4 | 4 | 4 | 3 | 1
0 | 8 | 10 | 8 | 8 | 3
1 | 1 | 2 | 1 | 16 | 2
Sliding Windows: Overlapping Time Intervals
Sliding windows create overlapping fixed-length intervals that advance by a specified hop distance. The _SlidingWindow class (lines 55‑61) manages these boundaries.
Window Assignment Mathematics
The _window_assignment_function() (lines 74‑124) builds an assign_windows() closure that calculates the first and last valid window indices (first_k, last_k) containing a given timestamp. It generates candidates via kth_stable_window(k) and filters to intervals where window_start <= timestamp < window_end.
Integration with Temporal Behavior
The _apply() method (lines 124‑165) executes this assignment function as a Pathway UDF, flattens the resulting window lists, and extracts _pw_window_start and _pw_window_end columns. When temporal behavior is specified, the method invokes utilities from temporal_behavior.py (lines 165‑203) to apply delay and cutoff parameters.
import pathway as pw
t = pw.debug.table_from_markdown(
'''
| instance | t
1 | 0 | 12
2 | 0 | 13
3 | 0 | 14
4 | 0 | 15
5 | 0 | 16
6 | 0 | 17
7 | 1 | 10
8 | 1 | 11
'''
)
result = t.windowby(
t.t,
window=pw.temporal.sliding(duration=10, hop=3),
instance=t.instance,
).reduce(
pw.this._pw_instance,
pw.this._pw_window_start,
pw.this._pw_window_end,
min_t=pw.reducers.min(pw.this.t),
max_t=pw.reducers.max(pw.this.t),
count=pw.reducers.count(),
)
pw.debug.compute_and_print(result, include_id=False)
Output:
_pw_instance | _pw_window_start | _pw_window_end | min_t | max_t | count
0 | 3 | 13 | 12 | 12 | 1
0 | 6 | 16 | 12 | 15 | 4
0 | 9 | 19 | 12 | 17 | 6
0 | 12 | 22 | 12 | 17 | 6
0 | 15 | 25 | 15 | 17 | 3
1 | 3 | 13 | 10 | 11 | 2
Tumbling Windows: Non-Overlapping Fixed Boundaries
Tumbling windows represent a specialized case of sliding windows where the hop distance equals the window duration, eliminating overlap. Rather than implementing a separate class, Pathway configures _SlidingWindow through the tumbling() helper function (lines 331‑337).
This helper instantiates _SlidingWindow with duration=None, hop=duration, and ratio=1, forcing each event to belong to exactly one discrete window. This architectural decision maximizes code reuse while providing a distinct API for non-overlapping analytics.
import pathway as pw
t = pw.debug.table_from_markdown(
'''
| instance | t
1 | 0 | 12
2 | 0 | 13
3 | 0 | 14
4 | 0 | 15
5 | 0 | 16
6 | 0 | 17
7 | 1 | 12
8 | 1 | 13
'''
)
result = t.windowby(
t.t,
window=pw.temporal.tumbling(duration=5),
instance=t.instance,
).reduce(
pw.this._pw_instance,
pw.this._pw_window_start,
pw.this._pw_window_end,
min_t=pw.reducers.min(pw.this.t),
max_t=pw.reducers.max(pw.this.t),
count=pw.reducers.count(),
)
pw.debug.compute_and_print(result, include_id=False)
Output:
_pw_instance | _pw_window_start | _pw_window_end | min_t | max_t | count
0 | 10 | 15 | 12 | 14 | 3
0 | 15 | 20 | 15 | 17 | 3
1 | 10 | 15 | 12 | 13 | 2
Temporal Behavior and Late Arrival Handling
Windowing operations in Pathway integrate with the temporal behavior system defined in python/pathway/stdlib/temporal/temporal_behavior.py. The CommonBehavior dataclass (lines 29‑34) encapsulates delay, cutoff, and keep_results parameters.
The apply_temporal_behavior() function (lines 101‑113) processes these configurations to buffer incoming events, freeze window state at cutoff times, or forget outdated data. This ensures that late-arriving events—those appearing after the theoretical window close—are either incorporated or discarded according to the specified policy.
The windowby() Entry Point
The public API surface for all windowing operations resides in the windowby() method (lines 560‑568 in _window.py). This method accepts a time expression, window instance, optional behavior configuration, and instance column, then delegates to the concrete window’s _apply() implementation. It serves as the unified interface that orchestrates the enrichment, flattening, and grouping pipeline described in the architecture section.
Summary
- Pathway provides three windowing primitives through
Table.windowby(): dynamic session windows via_SessionWindow, overlapping sliding windows via_SlidingWindow, and non-overlapping tumbling windows configured as a special case of sliding. - Session windows use predicate-based merging (
_merge()) and group representation pointers (_compute_group_repr()) to create variable-length intervals based on event proximity. - Sliding windows employ mathematical window assignment (
_window_assignment_function()) to calculate valid k-th windows for each timestamp, supporting overlapping intervals with configurable hop distances. - Tumbling windows reuse sliding logic with
ratio=1, eliminating overlap while maintaining the same flattening and aggregation infrastructure. - Temporal behavior handling (
temporal_behavior.py) provides configurable delay and cutoff mechanisms to manage late-arriving events across all window types. - All implementations reside in
python/pathway/stdlib/temporal/_window.py, with utility functions intemporal_behavior.pyandutils.py.
Frequently Asked Questions
What is the difference between sliding and tumbling windows in Pathway?
Sliding windows overlap and advance by a hop interval smaller than the window duration, causing events to appear in multiple windows. Tumbling windows are non-overlapping; each event belongs to exactly one window because the hop equals the duration. In the Pathway source code, both use _SlidingWindow, but tumbling() configures the constructor with ratio=1 to prevent overlap.
How does Pathway handle late-arriving data in windowing operations?
Pathway handles late arrivals through the temporal behavior system in temporal_behavior.py. The apply_temporal_behavior() function (lines 101‑113) implements delay parameters to buffer events and cutoff parameters to freeze window state, ensuring windows emit results only after the specified delay period has elapsed.
Can I use custom predicates for session windows?
Yes. The _SessionWindow class accepts a predicate parameter that defines custom grouping logic. The _merge() method (lines 70‑76) invokes this predicate via pw.apply_with_type to determine whether consecutive events belong to the same session, allowing complex business rules beyond simple time gaps.
Where are the windowing implementations located in the Pathway codebase?
The core windowing logic resides in python/pathway/stdlib/temporal/_window.py, which contains _SessionWindow, _SlidingWindow, and the windowby() entry point. Temporal behavior utilities are in python/pathway/stdlib/temporal/temporal_behavior.py, with supporting types in utils.py and join operations in _window_join.py.
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 →