andygrove opened a new pull request, #6684:
URL: https://github.com/apache/datafusion-comet/pull/6684

   ## Which issue does this PR close?
   
   Closes #6477.
   
   ## Rationale for this change
   
   DataFusion finds the `CURRENT ROW` bound of a `RANGE` window frame by 
comparing `ORDER BY` values with `ScalarValue::partial_cmp`, which puts a null 
array element or struct field above every other value. The sort puts it first, 
as Spark does. Once the bound search passes a row whose key holds such a null, 
the frame of every later row runs to the end of the partition. On `main`, 
`SUM(id) OVER (ORDER BY array(i))` over `(1, NULL), (2, 1), (3, 1), (4, 2)` 
returns 10 on every row, where Spark returns 1, 6, 6 and 10. `MAX`, `COUNT`, 
`LAST_VALUE`, a `DESC` key and `RANGE BETWEEN CURRENT ROW AND UNBOUNDED 
FOLLOWING` are wrong the same way, with no fallback and no error.
   
   This is the fallback-first fix. A native fix, a frame search that orders 
nested nulls the way the sort does, can follow.
   
   ## What changes are included in this PR?
   
   `CometWindowExec` already falls back for a `RANGE` frame that needs a bound 
search when DataFusion cannot compare an `ORDER BY` key at all 
(apache/datafusion#24937). It now also falls back, with its own reason, when an 
`ORDER BY` key's type can hold a null element or field, using 
`CometSortOrder.canHoldNestedNull`, which becomes public for this. The 
exemptions are the same as for that existing check: ranking functions and 
`ROWS` frames never search for a bound, `CUME_DIST`'s `RANGE` frame is never 
read, and a frame unbounded on both sides needs no search. All of them still 
matched Spark on `main` over these keys, and they stay native. So does a key 
whose type cannot hold a nested null, such as `array(coalesce(x, 0))`.
   
   Two fixtures used running sums over nullable nested keys to cover native 
`RANGE` frames, with the null row left out of the data. They now wrap the 
column in `coalesce` so they keep that native coverage, and 
`window_functions.sql` gains an `expect_fallback` for the nullable key.
   
   The strict floating-point check in `CometSortOrder` stays. It still covers 
sort keys under non-default null orders, which #6476 is about. Once that fix 
lands too, the strict-only check is redundant and can go, along with the 
`expect_fallback`s in `nested_float_order_keys_strict.sql`.
   
   The operator compatibility guide lists the new fallback.
   
   ## How are these changes tested?
   
   The new `windows/nested_null_range_frame.sql` fixture checks that these fall 
back with Spark's answers:
   - the default frame over an array key and over a struct key
   - a `DESC` key
   - `RANGE BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING`
   - `MAX` and `LAST_VALUE` over a struct of structs
   
   It also checks that `RANK`, `DENSE_RANK`, `ROW_NUMBER`, a `ROWS` frame, an 
unbounded `RANGE` frame, `CUME_DIST`, and keys that cannot hold a nested null 
stay native. With the new check disabled, the fixture fails on its first query 
with Comet's whole-partition sums.
   
   On Spark 4.1, these passed: every fixture under `windows/`, 
`CometWindowExecSuite`, `CometTopKSuite` and `CometFloatSemanticsSuite`. On 
Spark 3.5, the `windows/` fixtures and `CometWindowExecSuite` passed.
   


-- 
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