graph LR
shim_builder["shim_builder"]
stateful["stateful"]
flat_map_batch["flat_map_batch"]
join["join"]
fold_final["fold_final"]
window["window"]
build["build"]
fold_window["fold_window"]
shim_builder -- "generates stateful operators, utilizing" --> stateful
shim_builder -- "generates operators that utilize" --> flat_map_batch
shim_builder -- "generates operators for" --> fold_final
shim_builder -- "generates operators for" --> join
stateful -- "provides state to" --> join
stateful -- "provides state to" --> fold_final
build -- "configures" --> window
window -- "delegates aggregation to" --> fold_window
The bytewax.operators subsystem is central to defining and executing stream transformations. It comprises core primitives like flat_map_batch for efficient stateless processing and mechanisms for managing persistent state through the stateful component. The shim_builder acts as a versatile factory, responsible for constructing the internal logic of various complex operators, including those requiring state management (like joins and folds) and those that can leverage batch processing. The windowing components (build, window, fold_window) provide robust capabilities for time-based or event-based aggregations, with build configuring windowing strategies, window orchestrating the window lifecycle, and fold_window performing aggregations within these defined windows. This modular design allows for flexible and efficient stream processing, separating concerns between stateless transformations, state management, operator construction, and window-specific aggregations.
A versatile factory responsible for constructing the internal logic of various operators, including those for folds, collects, stateful flat maps, and joins. It abstracts the complexity of creating operator-specific processing units.
Related Classes/Methods: None
Manages persistent state across data elements, enabling stateful stream transformations. It provides the underlying mechanism for operators that require memory of past events or accumulated results.
Related Classes/Methods: None
A core primitive that many stateless transformations (e.g., map, filter, flatten) delegate to for efficient batch processing. It handles the application of a function to each element in a batch, potentially producing multiple output elements.
Related Classes/Methods:
Manages the complex operation of combining two or more data streams based on a common key, often involving state management to buffer and match elements.
Related Classes/Methods: None
Performs aggregations over an entire stream, producing a single final result once the stream is complete or a specific condition is met.
Related Classes/Methods: None
The central orchestrator within the windowing module that applies the defined windowing logic to incoming data. It manages the lifecycle of windows, including their creation, processing, and emission of results.
Related Classes/Methods:
The primary entry point for configuring and instantiating different windowing strategies (e.g., sliding, session, event-time, system-time).
Related Classes/Methods:
Performs aggregations specifically within the boundaries of defined windows, accumulating results per window.
Related Classes/Methods: