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

   ## Which issue does this PR close?
   
   Part of #5507. This fixes the part of it that is a regression in 1.1.0. 
#5507 stays open for the native sort order of nested floating-point values.
   
   Found by the 1.1.0 regression audit (#6399) and tracked in #6402.
   
   ## Rationale for this change
   
   #4870 made `WindowGroupLimit` native by default. Its rank cutoff finds 
`RANK` and `DENSE_RANK` ties by comparing the row-encoded `ORDER BY` key byte 
for byte. #5469 normalized scalar `FLOAT` and `DOUBLE` keys before that 
comparison, but a float nested in an array or struct keeps its raw bits. So 
`[-0.0]` and `[0.0]`, or two arrays holding different NaN encodings, get 
different ranks, and the cutoff drops rows that Spark keeps as ties.
   
   On 1.1.0-rc1, `RANK() OVER (ORDER BY array(v))` filtered to `rk <= 1`, over 
`v` = -0.0, 0.0 and 1.0, returns only the -0.0 row. Spark returns both tied 
rows. `DENSE_RANK() OVER (ORDER BY named_struct('x', v), id > 0)` does the 
same. In 1.0.0 `WindowGroupLimit` had no serde, so the limit and the window 
above it ran in Spark and were correct. The compatibility guide documented this 
as a known incompatibility rather than falling back.
   
   ## What changes are included in this PR?
   
   `CometWindowGroupLimitExec.convert` now falls back to Spark for `RANK` and 
`DENSE_RANK` when an `ORDER BY` key contains a `FLOAT` or `DOUBLE` nested at 
any depth, next to the existing check for collated keys. That restores the 
1.0.0 plan: the sort stays native, and the limit and the window run in Spark.
   
   - Scalar float keys are unaffected, because the native operator already 
normalizes them.
   - `ROW_NUMBER` keeps running natively. It never compares peers, so its 
cutoff is the first rows of the sorted input, which is also what Spark's limit 
took in 1.0.0 over the same native sort.
   - Partition keys need no check. Spark's `NormalizeFloatingNumbers` already 
normalizes nested floating-point partition keys before the limit is planned.
   
   The compatibility guide now lists the new fallback, and says that 
`ROW_NUMBER` over such a key still follows the native sort, whose order for 
nested floats can differ from Spark's (#5507).
   
   ## How are these changes tested?
   
   A new `CometWindowExecSuite` test runs `RANK` and `DENSE_RANK` with `rnk <= 
1` over `array(f)`, `array(d)` and `named_struct('x', d), id > 0`, with `-0.0`, 
`0.0` and `1.0` in a Parquet table. Each query has to match Spark, return both 
tied rows, and report the new fallback reason. The test also checks that a 
partitioned `ROW_NUMBER` over `array(d)` stays native and matches Spark, and 
that `RANK` partitioned by `array(d)` matches Spark.
   
   Without the serde change the new test fails with a result mismatch. With it, 
the whole of `CometWindowExecSuite` passes on the Spark 4.1 profile, and the 
`window group limit` tests pass on Spark 3.5.
   


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