apache / apache/datafusion-comet

Natively support time-window grouping expressions: window, session_window, window_time

Open
#4,553 2 comments 0 reactions 0 assignees View on GitHub
area:expressions enhancement expression priority:medium temporal expressions
Dominant language
Scala
Stars
1.3k
Forks
373
Avg merge
2d 4h
Merged PRs (30d)
198

Description

## Background

Spark's time-window grouping expressions currently fall back to Spark in Comet:

- `window(timeColumn, windowDuration, [slideDuration], [startTime])` (`TimeWindow`) - tumbling/sliding windows, common in batch aggregation (`GROUP BY window(ts, '1 hour')`), not just streaming.
- `session_window(timeColumn, gapDuration)` (`SessionWindow`) - session windows.
- `window_time(window)` (`WindowTime`) - extracts the event time from a window column.

Comet has no serde for `TimeWindow` / `SessionWindow` / `WindowTime` today, so any query using them falls back.

## Notes

These are not plain scalar functions: Spark's analyzer (`TimeWindowing` / `SessionWindowing` rules) rewrites `window()` / `session_window()` into an `Expand` plus grouping on a computed window struct. Native support would need to handle the rewritten form (struct construction and the window-boundary arithmetic) and the grouping that follows.

`window()` over a fixed duration is the most commonly used of the three and would be the natural starting point.

## Acceptance criteria

- `window`, `session_window`, and `window_time` execute natively in Comet and match Spark.
- Add SQL file test coverage under `expressions/datetime/`.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.