Skip to content

fix: fall back to Spark for RANK and DENSE_RANK limits over nested floating-point keys - #6468

Merged
andygrove merged 1 commit into
apache:mainfrom
andygrove:fix/wgl-nested-float-fallback-5507
Sep 30, 2026
Merged

andygrove merged 1 commit into
apache:mainfrom
andygrove:fix/wgl-nested-float-fallback-5507

Conversation

@andygrove

Copy link
Copy Markdown
Member

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.

…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.
@andygrove andygrove added bug Something isn't working backport-1.1 Candidate for backporting to 1.1 release branch labels Sep 30, 2026
@andygrove andygrove mentioned this pull request Sep 30, 2026
9 of 20 tasks
@andygrove
andygrove marked this pull request as ready for review September 30, 2026 17:22

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Native WindowGroupLimit distinguished nested floating-point encodings that Spark treats as peers, dropping tied rows at rank cutoffs.
  • Design approach: Fall back to Spark for RANK and DENSE_RANK when 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 WindowGroupLimit implementations 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.containsType keeps the change small without adding another abstraction. ROW_NUMBER and 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.

@andygrove
andygrove added this pull request to the merge queue Sep 30, 2026
@andygrove

Copy link
Copy Markdown
Member Author

Thanks @sunchao

Merged via the queue into apache:main with commit 300c748 Sep 30, 2026
40 checks passed
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 1, 2026
…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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport-1.1 Candidate for backporting to 1.1 release branch bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants