Repository navigation
[SPARK-59905][SQL][4.2] Fix storage-partitioned join sorting partitions the wrong way for an ORDER BY over a join-key expression - #59204
Closed
peter-toth wants to merge 1 commit into
Conversation
…ns the wrong way for an ORDER BY over a join-key expression
HyukjinKwon
approved these changes
Oct 2, 2026
dongjoon-hyun
approved these changes
Oct 2, 2026
dongjoon-hyun
left a comment
Member
There was a problem hiding this comment.
+1, LGTM. Thank you, @peter-toth .
dongjoon-hyun
pushed a commit
that referenced
this pull request
Oct 2, 2026
…ns the wrong way for an ORDER BY over a join-key expression ### Differences from #59189 This is the branch-4.2 backport of #59189. The description below is the original one. On this branch: - The change to `EnsureRequirements` does not cherry-pick cleanly, since the `OrderedDistribution` arm sits one level deeper here. The code is the same as on master, re-indented, and its comment is re-wrapped. - The two comment changes in `reorderJoinKeysRecursively` are left out. Branch-4.2 has no such comment there. - The test is the same as on master. Without the change, it fails on branch-4.2 the same way as on master. - Ran `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `ProjectedOrderingAndPartitioningSuite` and `PlannerSuite` on branch-4.2, 306 tests in all, plus `dev/lint-scala`. ### What changes were proposed in this pull request? When `EnsureRequirements` serves an `ORDER BY` by putting a keyed child's partitions in key order, it now compares the key rows position by position. Sort order `i` reads field `i` of the key row. It used to bind the ordering to the references of the partition expressions instead. The binding takes the types the key rows hold (`keyDataTypes`), and an assertion pins the precondition it relies on next to it. ### Why are the changes needed? A key row holds the values of the partition expressions. Binding the ordering to their references is right only when each expression is a bare column, the only thing a scan reports that an `ORDER BY` can name. A join can report the other side's join key instead. A one-side shuffle does, and so does an inner broadcast hash join, whose `expandOutputPartitioning` reports the streamed side's layout over the build side's join key. With `spark.sql.sources.v2.bucketing.shuffle.enabled` and `spark.sql.sources.v2.bucketing.sorting.enabled` on, this query returns `b` as 0 to 6 instead of 6 to 0. `ident(id)` is identity-partitioned and holds 94 to 100, and `plain(b)` holds 0 to 6: ```sql SELECT p.b FROM ident i JOIN plain p ON i.id = 100 - p.b ORDER BY 100 - p.b ``` The same happens with only `spark.sql.sources.v2.bucketing.sorting.enabled` when `plain` is broadcast. 1. The join shuffles `plain` onto `ident`'s layout and reports `100 - b` as `plain`'s partition expression. Its key rows hold 94 to 100. 2. `KeyedPartitioning.keysSatisfy` admits that partitioning for `ORDER BY 100 - p.b`, since its expression is the ordering's. So no range shuffle is planned, and the global sort is dropped as redundant. 3. `EnsureRequirements` binds `100 - b` to `b` and evaluates it on the key rows, so it computes `100 - (100 - b)`. The partitions end up in the reverse order. Two other join keys fail planning instead: - For a join on `i.id = p.b + p.c`, `ORDER BY p.b + p.c` binds two references to a key row with one field. Planning fails with an `ArrayIndexOutOfBoundsException`. - For a join on `i.id = p.s.a`, `ORDER BY p.s.a` binds the struct `s` to the long key and reads a field of it. Planning fails with a `ClassCastException`. `keysSatisfy` admits a partitioning here only when its expressions are the ordering's, position by position (`OrderedDistribution.areAllClusterKeysMatched`). So reading field `i` for sort order `i` is exact. ### Does this PR introduce _any_ user-facing change? Yes. Such a query returns its rows in the right order, and the two other forms no longer fail planning. They all need `spark.sql.sources.v2.bucketing.sorting.enabled`, which is off by default. A broadcast join needs nothing else, and a one-side shuffle also needs `spark.sql.sources.v2.bucketing.shuffle.enabled`. No range shuffle is added. ### How was this patch tested? New test in `KeyGroupedPartitioningSuite`, "SPARK-59905: ORDER BY a join-key expression a join reports as a partition expression". It runs five queries: - `ORDER BY 100 - p.b`, ascending and descending; - `ORDER BY p.s.a`; - `ORDER BY 100 - p.b, p.c DESC` over a two-key identity side, which needs a mixed-direction `GroupPartitionsExec`; - `ORDER BY p.b + p.c`. Each runs through a one-side shuffle and through a broadcast join, with AQE on and off. The `p.b + p.c` one runs with AQE off only, since AQE's plan validation fails on it for SPARK-59901, which #59165 fixes. Each run checks the order of the rows, the number of shuffles, and whether a `GroupPartitionsExec` sorts the partitions. All 18 runs fail on master: - the `100 - p.b` and two-key queries return their rows in the wrong order; - the `p.s.a` one fails with the `ClassCastException`; - the `p.b + p.c` one fails with the `ArrayIndexOutOfBoundsException`. Also ran the `KeyGroupedPartitioning*` suites, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `ProjectedOrderingAndPartitioningSuite` and `PlannerSuite`, 427 tests in all, plus `dev/lint-scala`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5.5) Closes #59204 from peter-toth/SPARK-59905-order-by-expression-key-4.2. Authored-by: Peter Toth <peter.toth@gmail.com> Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
Member
|
Merge Summary:
Posted by |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Differences from #59189
This is the branch-4.2 backport of #59189. The description below is the original one. On this branch:
EnsureRequirementsdoes not cherry-pick cleanly, since theOrderedDistributionarm sits one level deeper here. The code is the same as on master, re-indented, and its comment is re-wrapped.reorderJoinKeysRecursivelyare left out. Branch-4.2 has no such comment there.KeyGroupedPartitioningSuite,EnsureRequirementsSuite,GroupPartitionsExecSuite,ProjectedOrderingAndPartitioningSuiteandPlannerSuiteon branch-4.2, 306 tests in all, plusdev/lint-scala.What changes were proposed in this pull request?
When
EnsureRequirementsserves anORDER BYby putting a keyed child's partitions in key order, it now compares the key rows position by position. Sort orderireads fieldiof the key row. It used to bind the ordering to the references of the partition expressions instead. The binding takes the types the key rows hold (keyDataTypes), and an assertion pins the precondition it relies on next to it.Why are the changes needed?
A key row holds the values of the partition expressions. Binding the ordering to their references is right only when each expression is a bare column, the only thing a scan reports that an
ORDER BYcan name. A join can report the other side's join key instead. A one-side shuffle does, and so does an inner broadcast hash join, whoseexpandOutputPartitioningreports the streamed side's layout over the build side's join key.With
spark.sql.sources.v2.bucketing.shuffle.enabledandspark.sql.sources.v2.bucketing.sorting.enabledon, this query returnsbas 0 to 6 instead of 6 to 0.ident(id)is identity-partitioned and holds 94 to 100, andplain(b)holds 0 to 6:The same happens with only
spark.sql.sources.v2.bucketing.sorting.enabledwhenplainis broadcast.plainontoident's layout and reports100 - basplain's partition expression. Its key rows hold 94 to 100.KeyedPartitioning.keysSatisfyadmits that partitioning forORDER BY 100 - p.b, since its expression is the ordering's. So no range shuffle is planned, and the global sort is dropped as redundant.EnsureRequirementsbinds100 - btoband evaluates it on the key rows, so it computes100 - (100 - b). The partitions end up in the reverse order.Two other join keys fail planning instead:
i.id = p.b + p.c,ORDER BY p.b + p.cbinds two references to a key row with one field. Planning fails with anArrayIndexOutOfBoundsException.i.id = p.s.a,ORDER BY p.s.abinds the structsto the long key and reads a field of it. Planning fails with aClassCastException.keysSatisfyadmits a partitioning here only when its expressions are the ordering's, position by position (OrderedDistribution.areAllClusterKeysMatched). So reading fieldifor sort orderiis exact.Does this PR introduce any user-facing change?
Yes. Such a query returns its rows in the right order, and the two other forms no longer fail planning. They all need
spark.sql.sources.v2.bucketing.sorting.enabled, which is off by default. A broadcast join needs nothing else, and a one-side shuffle also needsspark.sql.sources.v2.bucketing.shuffle.enabled. No range shuffle is added.How was this patch tested?
New test in
KeyGroupedPartitioningSuite, "SPARK-59905: ORDER BY a join-key expression a join reports as a partition expression". It runs five queries:ORDER BY 100 - p.b, ascending and descending;ORDER BY p.s.a;ORDER BY 100 - p.b, p.c DESCover a two-key identity side, which needs a mixed-directionGroupPartitionsExec;ORDER BY p.b + p.c.Each runs through a one-side shuffle and through a broadcast join, with AQE on and off. The
p.b + p.cone runs with AQE off only, since AQE's plan validation fails on it for SPARK-59901, which #59165 fixes. Each run checks the order of the rows, the number of shuffles, and whether aGroupPartitionsExecsorts the partitions.All 18 runs fail on master:
100 - p.band two-key queries return their rows in the wrong order;p.s.aone fails with theClassCastException;p.b + p.cone fails with theArrayIndexOutOfBoundsException.Also ran the
KeyGroupedPartitioning*suites,EnsureRequirementsSuite,GroupPartitionsExecSuite,ProjectedOrderingAndPartitioningSuiteandPlannerSuite, 427 tests in all, plusdev/lint-scala.Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5.5)