Skip to content

[SPARK-59905][SQL][4.0] Fix storage-partitioned join sorting partitions the wrong way for an ORDER BY over a join-key expression - #59206

Closed
peter-toth wants to merge 2 commits into
apache:branch-4.0from
peter-toth:SPARK-59905-order-by-expression-key-4.0
Closed

peter-toth wants to merge 2 commits into
apache:branch-4.0from
peter-toth:SPARK-59905-order-by-expression-key-4.0

Conversation

@peter-toth

@peter-toth peter-toth commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

Differences from #59189

This is the branch-4.0 backport of #59189. The description below is the original one. This branch has different code, so the change and the test are tailored:

  • Master fixes the binding in EnsureRequirements.resolveChild, which this branch does not have. Here the same binding lives in EnsureRequirements.ensureOrdering, which sorts a KeyGroupedPartitioning's partition values for an ORDER BY. It bound the ordering to the leaves of the partition expressions, cast to Attribute.
  • The change binds sort order i to field i of the partition value there, with the same assertion. The type comes from the sort order's own expression, since this branch has no keyDataTypes.
  • Only the broadcast route is affected here. The one-side shuffle route already returns the right order. With only spark.sql.sources.v2.bucketing.sorting.enabled, SELECT /*+ BROADCAST(p) */ p.b FROM neg i JOIN plain p ON i.id = -p.b ORDER BY -p.b returns b in the wrong order, where neg(id) is identity-partitioned and holds -6 to 0.
  • A key with a literal fails planning here instead of returning the wrong order. For i.id = 100 - p.b, the leaf 100 fails the cast to Attribute, and planning fails with a ClassCastException.
  • The test is the master one with three changes:
    • it adds ORDER BY -p.b, ascending and descending, over neg, the shape that returns the wrong order here;
    • it drops the GroupPartitionsExec check, since this branch has no such node;
    • through a one-side shuffle it expects 2 shuffles instead of 1, since a range shuffle follows for the ORDER BY here.
  • Without the change, the broadcast runs of the -p.b queries return their rows in the wrong order on branch-4.0, and those of the 100 - p.b and two-key queries fail with the ClassCastException. The one-side shuffle runs, and the p.s.a and p.b + p.c queries, pass on this branch without the change. They are there to keep the diff to [SPARK-59905][SQL] Fix storage-partitioned join sorting partitions the wrong way for an ORDER BY over a join-key expression #59189 small.
  • Ran KeyGroupedPartitioningSuite, EnsureRequirementsSuite, ProjectedOrderingAndPartitioningSuite and PlannerSuite on branch-4.0, 189 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:

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)

…ns the wrong way for an ORDER BY over a join-key expression

@dongjoon-hyun dongjoon-hyun 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.

+1, LGTM. Thank you, @peter-toth .

@dongjoon-hyun

Copy link
Copy Markdown
Member

cc @uros-b because this is a correctness bug.

For the record, [VOTE] Release Spark 4.0.5 (RC1) started yesterday.

@peter-toth

peter-toth commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor Author

I added baa1281 with the rest of the master/4.3/4.2 test: the one-side shuffle route, and the p.s.a and p.b + p.c queries. They do not fail on this branch, but I'd add them to keep the diff to #59189 small. The description is updated to match.

@uros-b uros-b 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.

Thank you folks! will re-cut 4.0.5

@peter-toth

peter-toth commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor Author

@uros-b, although this is a correctness issue, it is't a regression. The bug has been there since 4.0 so it is up to you if you would like to include it in the 4.0.5 release or not.

@uros-b

uros-b commented Oct 3, 2026

Copy link
Copy Markdown
Member

Thank you @peter-toth, please feel free to merge the fix in any case, and we can decide whether to block current RC and re-cut, or just keep this for the next maintenance release instead. Based on the facts, it seems like shipping 4.1.4 / 4.0.5 without this fix does not in fact make existing users worse, so I think this can wait for the next train too.

Hence, we should be able to proceed with current [VOTE] Release Spark 4.0.5 (RC1) at this time too. So, I would say - let's keep the vote open and see if there are any additional issues / blockers. If not, then we can ship as-is, and this fix can wait for 4.0.6. WDYT? @peter-toth @dongjoon-hyun @szehon-ho (reviewers of the original fix PR)

HyukjinKwon

This comment was marked as outdated.

peter-toth added a commit that referenced this pull request Oct 4, 2026
…ns the wrong way for an ORDER BY over a join-key expression

### Differences from #59189

This is the branch-4.0 backport of #59189. The description below is the original one. This branch has different code, so the change and the test are tailored:
- Master fixes the binding in `EnsureRequirements.resolveChild`, which this branch does not have. Here the same binding lives in `EnsureRequirements.ensureOrdering`, which sorts a `KeyGroupedPartitioning`'s partition values for an `ORDER BY`. It bound the ordering to the leaves of the partition expressions, cast to `Attribute`.
- The change binds sort order `i` to field `i` of the partition value there, with the same assertion. The type comes from the sort order's own expression, since this branch has no `keyDataTypes`.
- Only the broadcast route is affected here. The one-side shuffle route already returns the right order. With only `spark.sql.sources.v2.bucketing.sorting.enabled`, `SELECT /*+ BROADCAST(p) */ p.b FROM neg i JOIN plain p ON i.id = -p.b ORDER BY -p.b` returns `b` in the wrong order, where `neg(id)` is identity-partitioned and holds -6 to 0.
- A key with a literal fails planning here instead of returning the wrong order. For `i.id = 100 - p.b`, the leaf `100` fails the cast to `Attribute`, and planning fails with a `ClassCastException`.
- The test is the master one with three changes:
  - it adds `ORDER BY -p.b`, ascending and descending, over `neg`, the shape that returns the wrong order here;
  - it drops the `GroupPartitionsExec` check, since this branch has no such node;
  - through a one-side shuffle it expects 2 shuffles instead of 1, since a range shuffle follows for the `ORDER BY` here.
- Without the change, the broadcast runs of the `-p.b` queries return their rows in the wrong order on branch-4.0, and those of the `100 - p.b` and two-key queries fail with the `ClassCastException`. The one-side shuffle runs, and the `p.s.a` and `p.b + p.c` queries, pass on this branch without the change. They are there to keep the diff to #59189 small.
- Ran `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite` and `PlannerSuite` on branch-4.0, 189 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 #59206 from peter-toth/SPARK-59905-order-by-expression-key-4.0.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
@peter-toth peter-toth closed this Oct 4, 2026
@peter-toth

Copy link
Copy Markdown
Contributor Author

Merge Summary:

Posted by merge_spark_pr.py

@peter-toth

Copy link
Copy Markdown
Contributor Author

Thank you @dongjoon-hyun, @uros-b, @HyukjinKwon for the review.

@uros-b, IMO this doesn't block 4.0.5 RC1.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants