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]
