fix: fall back to Spark for RANK and DENSE_RANK limits over nested floating-point keys - #6468
Conversation
…oating-point keys The native WindowGroupLimit finds RANK and DENSE_RANK ties by comparing row-encoded ORDER BY keys byte for byte. Scalar float keys are normalized first, but a float nested in an array or struct keeps its raw bits, so [-0.0] and [0.0] got different ranks and the cutoff dropped rows that Spark keeps as ties. In 1.0.0 the limit ran in Spark, so this is a regression in 1.1.0. Fall back for RANK and DENSE_RANK when an order key contains a nested FLOAT or DOUBLE. ROW_NUMBER never compares peers and stays native, and Spark already normalizes floating-point partition keys. Part of apache#5507.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native
WindowGroupLimitdistinguished nested floating-point encodings that Spark treats as peers, dropping tied rows at rank cutoffs. - Design approach: Fall back to Spark for
RANKandDENSE_RANKwhen order keys contain nested floats or doubles. - Correctness / compatibility analysis: The recursive type check covers nested containers while excluding scalar floating-point keys. Spark source confirms the peer-comparison semantics and recursive partition-key normalization. Relevant
WindowGroupLimitimplementations are identical across Spark 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Spark 3.4 remains excluded through its existing shim. - Key design decisions: Reusing
SupportLevel.containsTypekeeps the change small without adding another abstraction.ROW_NUMBERand scalar-key execution remain unchanged. Affected queries intentionally incur Spark execution and row-conversion costs. - Implementation sketch: Check order-key types before serialization, record the fallback reason, and return
None. Regression tests verify tied rows, the fallback diagnostic, and unaffected execution paths. Documentation describes the new policy. - Behavioral changes worth calling out: The affected limit and parent window execute in Spark, while the sort can remain native. The inherited nested-sort limitation tracked by #5507 remains outside this fix.
- Suggested improvements: None meeting the requested threshold. No introduced P1/P2 issues found within this review.
Reviewed the entire three-file diff at c312b1c193a62065bff7350227ef99ce3ea67f22 against 9f68a4144fdefc70ec26f3d80b0e4eafe72a8df4. The PR was not a draft. Snapshot and live checks contained no existing reviews, comments, or review threads. Routed skills: review-comet-pr and review-comet-expression-pr.
Exact-head CI: 25 checks passed and 15 were skipped. Required Checks passed. The Spark 4.1 execution job passed 1,157 tests, including the new nested-floating-point regression test: https://github.com/apache/datafusion-comet/actions/runs/36747740186.
Validation: A disposable Spark 4.1.3 probe passed 24 rank-limit cases covering nested float/double signed zeros and NaNs, plus 12 peer-comparison and partition-normalization scenarios. This probe exercised Spark semantics, not Comet integration. No local Comet/native build or benchmark was run. Spark SQL compatibility jobs were skipped in exact-head CI, and non-default Spark runtime coverage was not independently established.
|
Thanks @sunchao |
…ested float rank-limit guard Strict floating-point mode fell back for every float nested in a sort or window key, before those keys were normalized. Removing that fallback also admitted keys that can hold a null element or field, which the native sort and RANGE window frames order differently from Spark (apache#6476, apache#6477). CometSortOrder now keeps the fallback for a key that has a float and whose type can hold a nested null, and every other key stays native. apache#6468 made RANK and DENSE_RANK limits fall back for nested float order keys, because the native operator compared them raw. The native planner now normalizes those keys, so drop that guard and make its test expect native execution. Split the nested float order key fixture: the existing one runs with the default settings, and a new one runs in strict mode and pins both halves of the policy.
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
WindowGroupLimitnative by default. Its rank cutoff findsRANKandDENSE_RANKties by comparing the row-encodedORDER BYkey byte for byte. #5469 normalized scalarFLOATandDOUBLEkeys 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 tork <= 1, overv= -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.0WindowGroupLimithad 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.convertnow falls back to Spark forRANKandDENSE_RANKwhen anORDER BYkey contains aFLOATorDOUBLEnested 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.ROW_NUMBERkeeps 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.NormalizeFloatingNumbersalready normalizes nested floating-point partition keys before the limit is planned.The compatibility guide now lists the new fallback, and says that
ROW_NUMBERover 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
CometWindowExecSuitetest runsRANKandDENSE_RANKwithrnk <= 1overarray(f),array(d)andnamed_struct('x', d), id > 0, with-0.0,0.0and1.0in 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 partitionedROW_NUMBERoverarray(d)stays native and matches Spark, and thatRANKpartitioned byarray(d)matches Spark.Without the serde change the new test fails with a result mismatch. With it, the whole of
CometWindowExecSuitepasses on the Spark 4.1 profile, and thewindow group limittests pass on Spark 3.5.