bhaskargurram-ai opened a new pull request, #25958:
URL: https://github.com/apache/datafusion/pull/25958

   ## Which issue does this PR close?
   
   - Closes #24885.
   
   ## Rationale for this change
   
   The `first_value` / `last_value` **aggregate** functions can be used as 
window functions through the DataFrame API. That is currently the only way to 
combine first/last value over a window with `FILTER (WHERE ...)`. But as soon 
as the frame start can move, for example `ROWS BETWEEN 1 PRECEDING AND CURRENT 
ROW`, or a reversed `UNBOUNDED PRECEDING .. CURRENT ROW` frame (see #24884), 
the query fails:
   
   ```
   This feature is not implemented: Aggregate can not be used as a sliding 
accumulator because
   `retract_batch` is not implemented: last_value(?table?.v) ORDER BY 
[?table?.t ASC NULLS LAST]
   ROWS BETWEEN 1 PRECEDING AND CURRENT ROW
   ```
   
   ## What changes are included in this PR?
   
   `FirstValue` and `LastValue` now implement 
`AggregateUDFImpl::create_sliding_accumulator`, the same way `min` / `max` 
provide a separate sliding accumulator. It returns a new 
`SlidingFirstLastValueAccumulator` that supports `retract_batch`. The regular 
accumulators and the groups accumulators are unchanged, so normal aggregation 
keeps its O(1) state.
   
   How it works: a sliding frame adds rows at its end (`update_batch`) and 
retracts rows from its start (`retract_batch`), in the order they were added. 
With `FILTER`, the same filtered rows are passed to both calls. The accumulator 
numbers the rows in that order and keeps only the rows that can still be the 
result: every row, or only the non-null rows with `IGNORE NULLS`.
   - `first_value` keeps all such rows in a `VecDeque`, and the front is the 
result.
   - `last_value` keeps only the newest such row. It stays the result until it 
is retracted itself, and at that point no candidate is left in the frame.
   - `retract_batch(n rows)` advances a "rows retracted" counter and pops the 
candidates numbered below it.
   
   Each row is pushed and popped at most once, so the cost is amortized O(1) 
per row, and memory is O(frame) for `first_value` and O(1) for `last_value`. An 
empty frame, or a frame without a qualifying row, evaluates to a typed NULL.
   
   Ordering requirements: window aggregates do not get an `ORDER BY` of their 
own. `create_window_expr` builds them without `order_by`, and the window passes 
only the arguments to the accumulator, so "first" and "last" mean first and 
last in frame order. This matches how the existing trivial accumulators behave 
in plain (ever-expanding) windows. If an `ORDER BY` were present anyway, 
`create_sliding_accumulator` falls back to the existing accumulator (no 
retraction), so we never return results whose ordering semantics are not 
supported.
   
   ## What is the testing strategy for this PR?
   
   - Unit tests in `functions-aggregate` (`first_last.rs`):
     - `sliding_first_last_value_retract` and 
`sliding_first_last_value_retract_ignore_nulls`: step-by-step frames with 
NULLs, empty frames and refilling.
     - `sliding_first_last_value_matches_frames`: slides frames of length 1 to 
5 over 40 values with NULLs, for both null modes, and compares every result 
with the first/last (non-null) value of the frame computed directly.
     - `sliding_first_last_value_merge` and 
`sliding_first_last_value_with_order_by_falls_back`.
   - An end-to-end test in `physical-plan` (`bounded_window_agg_exec.rs`), 
`first_last_value_aggregate_sliding_frame`. It runs both `BoundedWindowAggExec` 
and `WindowAggExec` over two input batches with `ROWS BETWEEN 1 PRECEDING AND 
CURRENT ROW`. It checks `first_value` / `last_value` with `FILTER (WHERE 
keep)`, including a kept row whose value is NULL as in the issue's example, and 
with `IGNORE NULLS`.
   - Without the new `create_sliding_accumulator` overrides, the new tests 
fail. The end-to-end test then fails with the `retract_batch` "not implemented" 
error. A variant that drops one row too many on retract also fails 
`sliding_first_last_value_matches_frames`.
   
   No sqllogictest is included: SQL always resolves `first_value` / 
`last_value` in an `OVER` clause to the window UDFs 
(`SqlToRel::find_window_func` special-cases these names), so the aggregate path 
cannot be reached from SQL today. Routing `first_value(...) FILTER (...) OVER 
(...)` to the aggregate, as the issue suggests as a follow-up, would make 
SQL-level tests possible.
   
   Commands run: `cargo test -p datafusion-functions-aggregate --lib` (254 
passed), `cargo test -p datafusion-physical-plan --lib 
first_last_value_aggregate_sliding_frame` (1 passed), `cargo clippy -p 
datafusion-functions-aggregate -p datafusion-physical-plan --all-targets -- -D 
warnings` (clean) and `cargo fmt --all -- --check`.
   
   Questions for reviewers:
   - Is a separate sliding accumulator, as for `min`/`max`, the preferred 
shape? The alternative is adding retraction to the existing accumulators, but 
they would then have to buffer rows during regular aggregation too.
   - I made `merge_batch` of the sliding accumulator accept the `[value, 
is_set]` state of the trivial accumulators. I don't think window evaluation 
ever merges sliding accumulators, so this is only for completeness.
   
   ## Are there any user-facing changes?
   
   `first_value` / `last_value` aggregate UDAFs can now be used as window 
functions over sliding and reversed frames (DataFrame API), including with 
`FILTER`. There are no API changes.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to