Skip to content

feat(reporting): derive TimeWindowEvent event fact rows with mapInArrow - #119

Merged
tombonfert merged 5 commits into
mainfrom
feature/time_window_event_mapInArrow
Oct 9, 2026
Merged

tombonfert merged 5 commits into
mainfrom
feature/time_window_event_mapInArrow

Conversation

@tombonfert

Copy link
Copy Markdown
Collaborator

Summary

TimeWindowEvent.determine_events built its event_instance_fact rows with the pandas UDF window_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-level mapInArrow that emits one row per window, built from numpy arrays in bounded Arrow batches.

Closes #118

Changes

  • New explode_windows(df, *, id_col, start_col, stop_col, windows) in time_window_expression.py:
    • handles all events in one pass;
    • keeps the container_id type;
    • flushes output batches at about 100k windows, never splitting one event's windows for a container.
  • tile_windows remains the single window implementation, so event_instance_ids still match the scoped aggregations. event_instance_id, event_id and the fact schema are unchanged.
  • TimeWindowEvent.determine_events now uses explode_windows, and window_intervals_udf is removed.
  • Docstring fixes, including the truncated TimeWindowEvent API docs page.

… 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.
@tombonfert
tombonfert requested a review from a team as a code owner October 8, 2026 14:51
@codecov

codecov Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 94.44444% with 2 lines in your changes missing coverage. Please review.
✅ Project coverage is 90.17%. Comparing base (a44c2a5) to head (bf9929c).

Files with missing lines Patch % Lines
src/impulse_reporting/events/time_window_event.py 94.44% 1 Missing and 1 partial ⚠️
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     
Flag Coverage Δ
query_engine 86.66% <ø> (+0.02%) ⬆️
reporting 94.78% <94.44%> (-0.02%) ⬇️

Flags with carried forward coverage won't be shown. Click here to find out more.

Files with missing lines Coverage Δ
...ine/analyze/query/events/time_window_expression.py 97.14% <ø> (+2.20%) ⬆️
src/impulse_reporting/events/time_window_event.py 97.26% <94.44%> (-2.74%) ⬇️
🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

…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.
@tombonfert
tombonfert merged commit e26b13a into main Oct 9, 2026
6 checks passed
@tombonfert
tombonfert deleted the feature/time_window_event_mapInArrow branch October 9, 2026 12:55
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

TimeWindowEvent: derive event_instance_fact rows with mapInArrow

1 participant