Repository navigation
feat(reporting): derive TimeWindowEvent event fact rows with mapInArrow - #119
Merged
Merged
Conversation
… TimeWindowEvent window generation Replace `window_intervals_udf` with a new `explode_windows` helper that uses `mapInArrow` to emit one row per window directly. This processes multiple window lengths in a single pass, preserves the container id type, and bounds Python worker memory by flushing output batches of ~100,000 windows between containers. Update `TimeWindowEvent` to use the new helper and adjust tests/docs accordingly.
…tale container_id rename Fix the cross-reference in `TimeWindowEvent.determine_events` from `explode_windows` to `tile_windows`, matching the shared window function actually used. Remove the leftover `withColumnRenamed(solver.config.container_id_col, "container_id")` call in `determine_events` since `windows_df` already carries the correct `container_id` column. Sync the API reference docs with the corrected docstring and line wrapping.
Change `_window_batches` to emit output batches after completing one event's windows within a container, rather than after all events of a container. This keeps batch sizes bounded by `batch_windows - 1 + max_windows` independent of the number of events per container, preventing unbounded memory growth when many events share a container. Update docstrings and unit tests to reflect the per-event flush behavior and the tighter memory bound.
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #119 +/- ##
==========================================
+ Coverage 90.13% 90.17% +0.03%
==========================================
Files 66 66
Lines 5982 6003 +21
Branches 750 755 +5
==========================================
+ Hits 5392 5413 +21
+ Misses 468 467 -1
- Partials 122 123 +1
Flags with carried forward coverage won't be shown. Click here to find out more.
🚀 New features to boost your workflow:
|
…es configurable Add an `output_cols` parameter to `explode_windows` so callers can choose the names of the event name, window start, and window end columns instead of hardcoding `event_name`, `start_ts`, and `end_ts`. Validate that `id_col` and `output_cols` are four distinct names. Update `TimeWindowEvent` to pass its fact-schema column names and thread them through `generate_event_instance_id_column` and `ReportEntityUtil.get_event_id_column`. Adjust unit tests to cover custom output columns and collisions.
…gine into TimeWindowEvent Relocate the `explode_windows` window-generation helper and its batching helpers from `time_window_expression.py` into `TimeWindowEvent` as private `_explode_windows`, `_window_batches`, and `_record_batch`. This keeps the reporting event-fact implementation alongside the event class and removes the query engine's dependency on Spark DataFrame/Arrow APIs for this path. Update `TimeWindowEvent.determine_events` to call the local helper, and adjust docstrings and API reference docs to reference `TimeWindowEvent.determine_events` instead of the removed public function. Move the corresponding unit tests from the query-engine module to the reporting event tests.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
TimeWindowEvent.determine_eventsbuilt itsevent_instance_factrows with the pandas UDFwindow_intervals_udf. On the pandas UDF output path every window becomes a Python object, so worker memory grew with the containers per Arrow batch times their windows. Fine windows over long recordings could exhaust the Python worker. This PR replaces the UDF with a container-levelmapInArrowthat emits one row per window, built from numpy arrays in bounded Arrow batches.Closes #118
Changes
explode_windows(df, *, id_col, start_col, stop_col, windows)intime_window_expression.py:container_idtype;tile_windowsremains the single window implementation, soevent_instance_ids still match the scoped aggregations.event_instance_id,event_idand the fact schema are unchanged.TimeWindowEvent.determine_eventsnow usesexplode_windows, andwindow_intervals_udfis removed.TimeWindowEventAPI docs page.