Skip to content

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

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

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

Conversation

@peter-toth

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

Copy link
Copy Markdown
Contributor

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

…e wrong way for an ORDER BY over a join-key expression
@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you, @peter-toth .

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

Thank you for the fix, @peter-toth. I reviewed the latest commit, c988695, and left 7 inline comments. The change looks correct to me: keysSatisfy admits a partitioning for an OrderedDistribution only through areAllClusterKeysMatched, so reading key-row field i for sort order i is exact, and I found no correctness issue in the diff. Here is a summary, ordered by priority.

Tests

  1. KeyGroupedPartitioningSuite L5464: The same bug is also reachable through a broadcast hash join with only sorting.enabled. No test covers that route, and the PR description says shuffle.enabled is required.
  2. KeyGroupedPartitioningSuite L5478: AQE is off for all four queries, but only p.b + p.c needs that (SPARK-59901).
  3. KeyGroupedPartitioningSuite L5482: No multi-key case over a one-side-shuffled member with mixed sort directions.

API and design

  1. EnsureRequirements L146: The positional ordering relies on the keysSatisfy gate in another module, with no local check.

Docs and comments

  1. EnsureRequirements L145: Two copies of the old "single-column invariant" comment remain at L536 and L543.

Cleanup

  1. EnsureRequirements L148: Take the type from keyDataTypes rather than order.child.dataType.
  2. EnsureRequirements L149: order.copy(...) would avoid restating the direction and null ordering.

Items 1 and 2 are the most important. Also, #59165 rewrites the same comment block at EnsureRequirements L139-146 to say that SPARK-59905 is not handled yet, so whichever PR lands second will need to reconcile it.

}
}

test("SPARK-59905: ORDER BY the join-key expression of a one-side shuffle") {

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.

The same bug is also reachable through a broadcast hash join, with only sorting.enabled. BroadcastHashJoinExec.expandOutputPartitioning (BroadcastHashJoinExec.scala:119-121) substitutes the build-side join key into the partitioning of the streamed side, so an inner broadcast join turns KeyedPartitioning([i.id]) into a collection that also holds KeyedPartitioning([100 - p.b]) over the same layout. With these tables, spark.sql.sources.v2.bucketing.sorting.enabled=true and shuffle.enabled left at its default (false), this query reaches the same OrderedDistribution arm:

SELECT /*+ BROADCAST(p) */ p.b FROM testcat.ns.ident i JOIN testcat.ns.plain p ON i.id = 100 - p.b ORDER BY 100 - p.b

On master it should return b as 0 to 6 instead of 6 to 0, and the p.b + p.c and p.s.a forms should fail the same way as in this test. I traced this by reading the code and did not run it. This PR fixes that route too, but no test covers it, and the PR description says all the cases need spark.sql.sources.v2.bucketing.shuffle.enabled. Could you add a broadcast variant of these queries and update the description?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks, confirmed and done in 8f061e0. On master with only sorting.enabled, the broadcast form returns b as 0 to 6, and the p.s.a and p.b + p.c forms fail as in the shuffle case. The test now runs every query through both a one-side shuffle and a broadcast join, and the description no longer says the shuffle conf is needed.

The same expansion also reaches #59165's wrong result with every conf at its default (comment).

// `GroupPartitionsExec` to sort the partitions. AQE stays off, since its plan validation fails
// on `b + c` until SPARK-59901 is fixed.
withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",

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.

AQE is turned off for all four queries, but the stated reason only applies to p.b + p.c. The SPARK-59901 failure is the assert(refs.size == 1) in KeyedShuffleSpec.keyPositions (partitioning.scala:1954), which fires only for a key with two references. 100 - p.b and p.s.a have one reference each and map to an empty position set, so they should plan under AQE too. As written, the default AQE-on path is never tested for those three queries, e.g. the ValidateRequirements check in OptimizeSkewedJoin and the re-planning over the keyed ShuffleQueryStageExec. Could you run the three single-reference queries with AQE both off and on, and keep only p.b + p.c AQE-off until #59165 lands?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in 8f061e0. The single-column queries run with AQE on and off. Only p.b + p.c stays off until #59165 lands.

SQLConf.V2_BUCKETING_SHUFFLE_ENABLED.key -> "true",
SQLConf.V2_BUCKETING_SORTING_ENABLED.key -> "true") {
val ids = (94 to 100).map(i => Row(i.toLong))
Seq(

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.

Minor: Every case here has a single partition key. Two-key positional reading is already covered on a scan layout by SPARK-48655, but no test covers a multi-key member built by the one-side shuffle, where key-row field i has to line up with the i-th expression that KeyedShuffleSpec.createPartitioning emits and the direction differs per position. For example, ident2(x, y) partitioned by identity(x), identity(y) with keys (94, 1), (94, 2), (95, 1), joined ON i.x = 100 - p.b AND i.y = p.c and ordered by 100 - p.b, p.c DESC, which needs a mixed-direction GroupPartitionsExec. Could you add such a case?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Added in 8f061e0, with your ident2(x, y) shape ordered by 100 - p.b, p.c DESC. It needs the mixed-direction GroupPartitionsExec, and master returns its rows in the wrong order on both routes.

// partitioning here only when those expressions are the ordering's, position by
// position, so each sort order reads its position of the key row. The expression
// need not be a bare column, e.g. a shuffled side's join key `100 - b`.
val keyRowOrdering = RowOrdering.create(

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.

Minor: This ordering is only right because KeyedPartitioning.keysSatisfy admits the partitioning through o.areAllClusterKeysMatched(expressions) (partitioning.scala:963), in another module, and nothing here checks it. If that gate is ever relaxed, e.g. to admit an ordering prefix the way RangePartitioning.satisfies0 does, then KeyedPartitioning([a]) under ORDER BY a, b would read field 1 of a one-field key row, and a permuted match would sort the partitions in the wrong order without any error. Could you add assert(o.areAllClusterKeysMatched(satisfyingKeyedPartitioning.expressions)) here, or build this ordering in KeyedPartitioning next to keysSatisfy?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in 8f061e0, with the assertion.

// A key row holds the values of the partition expressions. `keysSatisfy` admits a
// partitioning here only when those expressions are the ordering's, position by
// position, so each sort order reads its position of the key row. The expression
// need not be a bare column, e.g. a shuffled side's join key `100 - b`.

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.

Minor: This new comment is right, but two copies of the comment it replaces are still in this file, at L536 and L543 in reorderJoinKeysRecursively: "The single-column invariant in KeyedPartitioning.supportsExpressions guarantees one attribute per partition expression." supportsExpressions is only checked where a scan reports its partitioning (DataSourceV2ScanExecBase.scala:105). For the p.b + p.c join in this test, kp.expressions.flatMap(_.references) there gives two attributes for one expression. The code is harmless (reorder gives up on a size mismatch), but this assumption is what caused SPARK-59905, and #59165 only rewords the third copy, at L1130. Could you update these two as well?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in 8f061e0. Both now say that a scan reports one column per expression, and that reorder gives up on a join key like b + c.

// need not be a bare column, e.g. a shuffled side's join key `100 - b`.
val keyRowOrdering = RowOrdering.create(
o.ordering.zipWithIndex.map { case (order, i) =>
SortOrder(BoundReference(i, order.child.dataType, nullable = true),

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.

Minor: The type here is order.child.dataType, which is the partition expression type, not the type the key rows are compared at. The KeyedPartitioning doc says expressionDataTypes "Need not be what the partitionKeys rows hold. See keyDataTypes", and that ShuffleExchangeExec "is the one reader that stays on expressionDataTypes" (partitioning.scala:704 and :723). Today the two differ only in struct field names and nullability, which the generated ordering ignores, so the result is right. But that depends on join type coercion and the reduce gates keeping them in step. Could you zip o.ordering with satisfyingKeyedPartitioning.keyDataTypes and use that type? The lengths match because areAllClusterKeysMatched uses corresponds.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in 8f061e0.

val keyRowOrdering = RowOrdering.create(
o.ordering.zipWithIndex.map { case (order, i) =>
SortOrder(BoundReference(i, order.child.dataType, nullable = true),
order.direction, order.nullOrdering, Seq.empty)

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.

Nit: order.copy(child = BoundReference(i, ..., nullable = true), sameOrderExpressions = Seq.empty) would say "the same sort order, pointed at key position i" without restating direction and nullOrdering. ShuffleExchangeExec.scala:393-395 uses that form for the RangePartitioner sample ordering. Keep nullable = true rather than copying order.nullable as that code does, since a key row can hold a null where the ORDER BY expression is not nullable. Could you use copy here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in 8f061e0.

@peter-toth

Copy link
Copy Markdown
Contributor Author

Thanks for the review, @dongjoon-hyun. 8f061e0 takes all seven items. I updated the description for the broadcast route, which does not need the shuffle conf.

@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 for addressing all the comments, @peter-toth.

@szehon-ho szehon-ho 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.

nice find!

dongjoon-hyun pushed a commit that referenced this pull request Oct 1, 2026
…e wrong way for an ORDER BY over a join-key expression

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

Closes #59189 from peter-toth/SPARK-59905-order-by-expression-key.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
(cherry picked from commit ec2a1bb)
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
@dongjoon-hyun

Copy link
Copy Markdown
Member

Merge Summary:

Posted by merge_spark_pr.py

@dongjoon-hyun

Copy link
Copy Markdown
Member

Could you make backporting PRs to older release branches, @peter-toth ? I guess we need at least at branch-4.3 and branch-4.2.

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.3 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 test is the same as on master. Without the change, it fails on branch-4.3 the same way as on master.
- Ran `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `ProjectedOrderingAndPartitioningSuite` and `PlannerSuite` on branch-4.3, 360 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 #59203 from peter-toth/SPARK-59905-order-by-expression-key-4.3.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
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>
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.1 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.1, 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.1, 210 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 #59205 from peter-toth/SPARK-59905-order-by-expression-key-4.1.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
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 added a commit that referenced this pull request Oct 5, 2026
…itioned join pairs a transform of a join-key expression with one of a column

### What changes were proposed in this pull request?

A storage-partitioned join now pairs a partition transform whose argument is not a bare column, i.e. is an expression or a struct field, only with a transform over the same argument. It also stops dropping such an argument where it used to:

- `KeyedShuffleSpec.isExpressionCompatible` pairs a transform with an argument that is not an `Attribute` only with a transform over the same argument shape, i.e. the argument with its one column replaced by a placeholder. So `bucket(4, b + 1)` pairs with `bucket(4, c + 1)` on `b = c`, but not with `bucket(4, x)`. Over the same shape the two compare as two transforms of columns do. Under `allowCompatibleTransforms`, `bucket(8, b + 1)` reduces onto `bucket(4, c + 1)`, since a reduce maps the partition keys and evaluates no argument. When an earlier join reduced the keys, two such transforms pair only when they were reduced through the same pairing. An identity side is never reduced onto such a transform. `TransformExpression` keeps its comparisons. It loses `withReference`, whose last callers now rebuild the transform with `copy`.
- `KeyedShuffleSpec.canCreatePartitioning` asks the same question, so such a layout is never the one other children are shuffled onto. `createPartitioning` replaces a transform's argument with the other child's cluster key, which drops whatever surrounds the key. It and the identity arm of `reducersBothWays` now rebuild a transform through one helper that asserts the argument is a column.
- `ShuffleSpecCollection.canCreatePartitioning` is now true when any member can, and `EnsureRequirements.pickCoPartitionTarget` offers and ranks only the members that can. So one member that cannot serve no longer rules out a usable sibling. That also covers a bare expression such as `b + 1`, which master already refused.
- `GroupPartitionsExec` rebuilds a reduced expression over each member's own argument when it reports the members of a collection. It used to re-target the column only.
- `KeyedShuffleSpec.keyPositions` maps an expression with more than one reference, e.g. `bucket(4, b + c)`, to no position instead of failing an assertion (SPARK-59901).
- `EnsureRequirements.createKeyedShuffleSpecs` counts only an expression over one column towards `spark.sql.requireAllClusterKeysForCoPartition`. An expression over two columns maps to no position, so a projection would drop it together with its columns.

### Why are the changes needed?

With `spark.sql.sources.v2.bucketing.shuffle.enabled` on, this query returns 0 rows instead of 7. `t1(id)` and `t3(x)` are partitioned by `bucket(4, ...)`, `plain(b)` is not partitioned, and each holds the values 0 to 7:

```sql
SELECT t1.id, p.b, t3.x FROM t1
JOIN plain p ON t1.id = p.b + 1
JOIN t3 ON p.b = t3.x
```

1. The first join shuffles `plain` onto `t1`'s partitioning. `KeyedShuffleSpec.createPartitioning` builds that partitioning from `plain`'s join key, so it is `bucket(4, b + 1)`. This join is correct.
2. The second join clusters `plain` on the bare `b`. `KeyedShuffleSpec.keyPositions` maps a partition expression to a cluster key through its reference, so `bucket(4, b + 1)` counts as a function of `b`, the counterpart of `x`.
3. `isSameFunction` compares only the function name and the bucket count. So `bucket(4, b + 1)` is the same as `t3`'s `bucket(4, x)`, and the join pairs the partitions as they stand. A row with `b = x` sits in bucket `(b + 1) % 4` on one side and in `x % 4` on the other.

The comparison ignores which column an argument is, and relies on `keyPositions` to pair the columns up. That is sound only when the argument is the column itself. `bucket(4, b + 1)` is a function of `b`, but not the same function of it as `bucket(4, x)` is of `x`. A scan never reports a transform of an expression or of a struct field, but two planner paths build one. A one-side shuffle does, under `spark.sql.sources.v2.bucketing.shuffle.enabled`. An inner broadcast hash join does too, with every conf at its default. `BroadcastHashJoinExec.expandOutputPartitioning` also reports the streamed side's layout over the build side's join key. The same problem shows up in more shapes, all measured on master:

- **A broadcast join.** With every conf at its default, `SELECT /*+ BROADCAST(p) */ ... FROM t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x` returns 0 of 7 rows the same way, through the expanded `bucket(4, b + 1)` layout.
- **The reducer path.** Under `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`, a `bucket(8, x)` third table loses the rows the same way.
- **A struct field.** A join key `p.s.a` gives `bucket(4, s.a)`, which `keyPositions` pairs through `s`. Joining two such sides on `s` pairs `bucket(4, s.a)` with `bucket(4, s.b)` and returns 0 of 8 rows. Shuffling another side onto such a layout builds `bucket(4, s)`, which fails on executors with a `ClassCastException`.
- **Shuffling onto the layout.** In `plain p JOIN t1 ON p.b + 1 = t1.id JOIN xy q ON t1.id = q.x AND p.b = q.y`, the second join shuffles `xy` onto the first join's `bucket(4, b + 1)` layout. That builds `bucket(4, y)` for `xy`, drops the `+ 1`, and returns 0 of 8 rows.
- **After a reduce.** When a later join reduces the first join's output from `bucket(4)` onto `bucket(2)`, `GroupPartitionsExec` reports the `bucket(4, b + 1)` member as `bucket(2, b)`. A join on `b` then pairs it with a `bucket(2, x)` scan, or shuffles another side onto it, and returns 0 of 7 rows. The same happens to the bare `b + 1` member of a side shuffled onto an identity-partitioned table. The struct form fails with the `ClassCastException`.
- **An identity side reduced onto it.** Under `allowCompatibleTransforms`, a later join can reduce an identity-partitioned side onto `bucket(4, b + 1)`. That evaluates `bucket(4, id + 1)` on the side's partition keys while planning. The keys come out right, but under ANSI a key of `Long.MaxValue` fails the query with `ARITHMETIC_OVERFLOW`, although the query never computes `id + 1`.
- **A GROUP BY after a reduce.** With `allowCompatibleTransforms`, `SELECT /*+ BROADCAST(p) */ p.b, max(p.c) FROM ident i JOIN plain2 p ON i.id = p.b + p.c JOIN bucket2 t2 ON i.id = t2.id GROUP BY p.b` returns 4 groups instead of 2. The reduce reports the `p.b + p.c` member as `bucket(2, p.b)`, which serves the `GROUP BY` as it stands.
- **A key over two columns.** A join key such as `p.b + p.c` gives `bucket(4, b + c)`, or the bare `b + c` over an identity-partitioned side. Any spec over it fails planning with an `AssertionError` in `keyPositions` (SPARK-59901). AQE validates the plan in `OptimizeSkewedJoin`, so a single join is enough.

An `ORDER BY` over such a partition expression had a related problem in another code path. That was SPARK-59905, fixed in #59189.

The one-side shuffle of transform expressions came with SPARK-48012, in 4.0.0. So did the broadcast expansion of a keyed layout, since SPARK-49205 made `KeyGroupedPartitioning` an expression.

### Does this PR introduce _any_ user-facing change?

Yes, it fixes the wrong results, the crash and the planning failures above. Such a query now shuffles the side instead of pairing it or shuffling onto it.

That costs shuffles in three shapes whose plan was already correct:
- A widening cast counts as part of the argument's shape. When `t1.id` is `BIGINT` and `p.b` is `INT`, the first join reports `bucket(4, cast(b as bigint))`. A later join `p.b = t3.x` with an `INT` bucketed `t3` then shuffles both sides. Before, a connector whose bucket function has one canonical name for both types paired the two sides as they stood. This happens with every conf at its default, through a broadcast join. Two sides that both carry the cast still pair.
- Under `allowCompatibleTransforms`, an identity-partitioned side is no longer reduced onto such a transform, so it is shuffled.
- With `spark.sql.sources.v2.bucketing.shuffle.enabled`, a join on two keys whose sides pair only through a transform of an expression gets one more shuffle. Take two broadcast joins, `t1 JOIN /*+ BROADCAST(p) */ p ON t1.id = p.b + 1` and the same over `t2` and `q`, joined on `l.id = r.c AND l.b = r.b`. Each layout covers only one join key, so the storage-partitioned join declines under `requireAllClusterKeysForCoPartition`. Master then keeps both sides, paired on `bucket(4, b + 1)`. This PR cannot shuffle onto that layout, so it shuffles one side onto `bucket(4, id)`. Both plans co-partition on one key only, which `requireAllClusterKeysForCoPartition` is meant to prevent. That is a separate, pre-existing question, filed as SPARK-59971.

One plan gets better. A collection with a bare expression member, e.g. the `b + 1` of a side shuffled onto an identity-partitioned table, used to be refused as a layout as a whole. Now its other members can serve, so `ident i JOIN plain p ON i.id = p.b + 1 JOIN xy q ON i.id = q.x AND p.b = q.y` takes 2 shuffles instead of 3.

This PR keeps the fix small, since it has to reach the maintenance branches. SPARK-59900 is the follow-up for the rest: shuffling another side onto such a layout, and reducing an identity side onto it.

### How was this patch tested?

New tests. All but two of them fail on master:
- `ShuffleSpecSuite`, "SPARK-59887: a spec over a transform of an expression pairs only with the same shape". For an expression and a struct field argument, it covers the same shape, another shape and a column, both ways and with itself, the reduce in both directions, the two sides of one reduce, two sides reduced through different pairings, an identity side, `canCreatePartitioning`, and a collection with a usable sibling.
- `ShuffleSpecSuite`, "SPARK-59901: an expression over two columns maps to no cluster key".
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side shuffled onto a transform of an expression is not paired as it is". It runs the query above against a `bucket(4)` third table, and with `allowCompatibleTransforms` against a `bucket(8)` one that the join reduces, and asserts two shuffles.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side is not shuffled onto a transform of an expression". It runs both join orders and asserts that `xy` is shuffled onto `bucket(4, x)` and never onto a bucket of `y`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a broadcast join's layout over a join key is not paired as it is". It runs with every conf at its default and checks that the first join is a broadcast join.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a layout over a transform of a struct field is not paired nor shuffled onto". Master returns 0 of 8 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps what each member's keys are computed from". It covers a bucketed and an identity-partitioned first table, a bucketed and an unpartitioned third table, and the struct form. Master returns 0 of 7 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: an identity side is not reduced onto a transform of an expression". It runs with `p.b + 1` over a `Long.MaxValue` key and with `-p.b` over a `Long.MinValue` key. Master fails with the `ARITHMETIC_OVERFLOW` for both.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps a member over two columns". It runs the `GROUP BY` above. Master returns 4 groups instead of 2.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns is not paired through one of them". Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns does not cover its columns". It joins two broadcast joins that report `[b + c, d]` on `(b, c, d)` under `allowKeysSubsetOfPartitionKeys`, and asserts that the join is not paired on `d` alone. Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: two layouts over the same shape of a join key pair as two columns do". It runs the cast case above with every conf at its default. It also runs two one-side shuffles onto `bucket(4, b + 1)` and onto `bucket(4)` or `bucket(8)` of `q.b + 1`, where `allowCompatibleTransforms` reduces the `bucket(8)` one. It asserts that the join adds no shuffle. It passes on master, which pairs and reduces these as well.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: the same shape reduced together pairs only through the same pairing". Two legs each shuffle onto `bucket(12, b + 1)`. The first leg reduces with `bucket(8)`, so its keys are buckets of 4. The second leg reduces with `bucket(8)` too, with `bucket(18)`, or not at all, so its keys are buckets of 4, 6 or 12. It asserts 2 shuffles for the first and 4 for the other two. Master returns 0 of 23 rows for the third.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a member that cannot serve does not rank its collection". It pins the ranking: `q` keeps its layout with 3 partitions, as on master, where the collection is refused as a whole. Without the member filter, the collection would rank first with 5 partitions, and `q` would be shuffled onto 2.

The SPARK-59905 test from #59189 now also runs its `p.b + p.c` query with AQE on. Master fails that run with the `AssertionError`.

The three expression tests each run with the key `p.b + 1` and with `-p.b`. Master returns 0 of 8 rows for `p.b + 1` and 4 of 8 for `-p.b`. The `-p.b` key is the shape that also fails on 4.2, 4.1 and 4.0. On those branches `b + 1` already works, since its literal leaf blocks the pairing.

Each change was removed on its own:
- without the guard in `isExpressionCompatible`, the one-side shuffle and broadcast pairing tests, the struct field test and the reduce test return wrong rows, and the identity tests fail on the assertion in `rebuiltOver`;
- without the shape pairing, the same-shape and same-pairing tests fail on the shuffle count;
- with a reduce of such a pair refused, the same-shape test fails on the shuffle count;
- with every pair of reduced keys over an expression refused, the same-pairing test fails on the shuffle count;
- with the reduced marks ignored, or only checked on both sides, the same-pairing test returns 7 of 23 rows;
- without the `canCreatePartitioning` clause, the one-side shuffle pairing, `xy`, struct field, reduce and same-pairing tests fail, all on the assertion in `rebuiltOver`;
- with the `forall` back in `ShuffleSpecCollection.canCreatePartitioning`, the `xy` tests fail on the shuffle count;
- without the member filter in `pickCoPartitionTarget`, 15 tests fail in `KeyGroupedPartitioningSuite` and `EnsureRequirementsSuite`. The filter now also does what the collection-level check before it did, so 8 of them are older tests. Of this PR's tests, 7 fail, 6 of them on the assertion in `rebuiltOver`. A member over the same shape pairs with itself, so only the filter keeps it from being the layout;
- without the `GroupPartitionsExec` change, the reduce test returns 0 of 7 rows, and the two-column `GROUP BY` returns 4 groups;
- without the `requireAllClusterKeysForCoPartition` change, the coverage test pairs on `d` alone;
- with the `keyPositions` assertion back, both SPARK-59901 tests in `KeyGroupedPartitioningSuite` fail.

The first `ShuffleSpecSuite` test also fails for each change to `isExpressionCompatible` above.

The `xy` tests pass with the guard removed, since the `canCreatePartitioning` clause and the member filter keep `xy` off that layout.

Also ran `TransformExpressionSuite`, `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `WriteDistributionAndOrderingSuite`, `PlannerSuite`, `ProjectedOrderingAndPartitioningSuite`, `GroupPartitionsExecSuite`, `EnsureRequirementsSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite`, 1955 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 #59165 from peter-toth/SPARK-59887-spj-expression-key-identity.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
peter-toth added a commit that referenced this pull request Oct 5, 2026
…itioned join pairs a transform of a join-key expression with one of a column

### What changes were proposed in this pull request?

A storage-partitioned join now pairs a partition transform whose argument is not a bare column, i.e. is an expression or a struct field, only with a transform over the same argument. It also stops dropping such an argument where it used to:

- `KeyedShuffleSpec.isExpressionCompatible` pairs a transform with an argument that is not an `Attribute` only with a transform over the same argument shape, i.e. the argument with its one column replaced by a placeholder. So `bucket(4, b + 1)` pairs with `bucket(4, c + 1)` on `b = c`, but not with `bucket(4, x)`. Over the same shape the two compare as two transforms of columns do. Under `allowCompatibleTransforms`, `bucket(8, b + 1)` reduces onto `bucket(4, c + 1)`, since a reduce maps the partition keys and evaluates no argument. When an earlier join reduced the keys, two such transforms pair only when they were reduced through the same pairing. An identity side is never reduced onto such a transform. `TransformExpression` keeps its comparisons. It loses `withReference`, whose last callers now rebuild the transform with `copy`.
- `KeyedShuffleSpec.canCreatePartitioning` asks the same question, so such a layout is never the one other children are shuffled onto. `createPartitioning` replaces a transform's argument with the other child's cluster key, which drops whatever surrounds the key. It and the identity arm of `reducersBothWays` now rebuild a transform through one helper that asserts the argument is a column.
- `ShuffleSpecCollection.canCreatePartitioning` is now true when any member can, and `EnsureRequirements.pickCoPartitionTarget` offers and ranks only the members that can. So one member that cannot serve no longer rules out a usable sibling. That also covers a bare expression such as `b + 1`, which master already refused.
- `GroupPartitionsExec` rebuilds a reduced expression over each member's own argument when it reports the members of a collection. It used to re-target the column only.
- `KeyedShuffleSpec.keyPositions` maps an expression with more than one reference, e.g. `bucket(4, b + c)`, to no position instead of failing an assertion (SPARK-59901).
- `EnsureRequirements.createKeyedShuffleSpecs` counts only an expression over one column towards `spark.sql.requireAllClusterKeysForCoPartition`. An expression over two columns maps to no position, so a projection would drop it together with its columns.

### Why are the changes needed?

With `spark.sql.sources.v2.bucketing.shuffle.enabled` on, this query returns 0 rows instead of 7. `t1(id)` and `t3(x)` are partitioned by `bucket(4, ...)`, `plain(b)` is not partitioned, and each holds the values 0 to 7:

```sql
SELECT t1.id, p.b, t3.x FROM t1
JOIN plain p ON t1.id = p.b + 1
JOIN t3 ON p.b = t3.x
```

1. The first join shuffles `plain` onto `t1`'s partitioning. `KeyedShuffleSpec.createPartitioning` builds that partitioning from `plain`'s join key, so it is `bucket(4, b + 1)`. This join is correct.
2. The second join clusters `plain` on the bare `b`. `KeyedShuffleSpec.keyPositions` maps a partition expression to a cluster key through its reference, so `bucket(4, b + 1)` counts as a function of `b`, the counterpart of `x`.
3. `isSameFunction` compares only the function name and the bucket count. So `bucket(4, b + 1)` is the same as `t3`'s `bucket(4, x)`, and the join pairs the partitions as they stand. A row with `b = x` sits in bucket `(b + 1) % 4` on one side and in `x % 4` on the other.

The comparison ignores which column an argument is, and relies on `keyPositions` to pair the columns up. That is sound only when the argument is the column itself. `bucket(4, b + 1)` is a function of `b`, but not the same function of it as `bucket(4, x)` is of `x`. A scan never reports a transform of an expression or of a struct field, but two planner paths build one. A one-side shuffle does, under `spark.sql.sources.v2.bucketing.shuffle.enabled`. An inner broadcast hash join does too, with every conf at its default. `BroadcastHashJoinExec.expandOutputPartitioning` also reports the streamed side's layout over the build side's join key. The same problem shows up in more shapes, all measured on master:

- **A broadcast join.** With every conf at its default, `SELECT /*+ BROADCAST(p) */ ... FROM t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x` returns 0 of 7 rows the same way, through the expanded `bucket(4, b + 1)` layout.
- **The reducer path.** Under `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`, a `bucket(8, x)` third table loses the rows the same way.
- **A struct field.** A join key `p.s.a` gives `bucket(4, s.a)`, which `keyPositions` pairs through `s`. Joining two such sides on `s` pairs `bucket(4, s.a)` with `bucket(4, s.b)` and returns 0 of 8 rows. Shuffling another side onto such a layout builds `bucket(4, s)`, which fails on executors with a `ClassCastException`.
- **Shuffling onto the layout.** In `plain p JOIN t1 ON p.b + 1 = t1.id JOIN xy q ON t1.id = q.x AND p.b = q.y`, the second join shuffles `xy` onto the first join's `bucket(4, b + 1)` layout. That builds `bucket(4, y)` for `xy`, drops the `+ 1`, and returns 0 of 8 rows.
- **After a reduce.** When a later join reduces the first join's output from `bucket(4)` onto `bucket(2)`, `GroupPartitionsExec` reports the `bucket(4, b + 1)` member as `bucket(2, b)`. A join on `b` then pairs it with a `bucket(2, x)` scan, or shuffles another side onto it, and returns 0 of 7 rows. The same happens to the bare `b + 1` member of a side shuffled onto an identity-partitioned table. The struct form fails with the `ClassCastException`.
- **An identity side reduced onto it.** Under `allowCompatibleTransforms`, a later join can reduce an identity-partitioned side onto `bucket(4, b + 1)`. That evaluates `bucket(4, id + 1)` on the side's partition keys while planning. The keys come out right, but under ANSI a key of `Long.MaxValue` fails the query with `ARITHMETIC_OVERFLOW`, although the query never computes `id + 1`.
- **A GROUP BY after a reduce.** With `allowCompatibleTransforms`, `SELECT /*+ BROADCAST(p) */ p.b, max(p.c) FROM ident i JOIN plain2 p ON i.id = p.b + p.c JOIN bucket2 t2 ON i.id = t2.id GROUP BY p.b` returns 4 groups instead of 2. The reduce reports the `p.b + p.c` member as `bucket(2, p.b)`, which serves the `GROUP BY` as it stands.
- **A key over two columns.** A join key such as `p.b + p.c` gives `bucket(4, b + c)`, or the bare `b + c` over an identity-partitioned side. Any spec over it fails planning with an `AssertionError` in `keyPositions` (SPARK-59901). AQE validates the plan in `OptimizeSkewedJoin`, so a single join is enough.

An `ORDER BY` over such a partition expression had a related problem in another code path. That was SPARK-59905, fixed in #59189.

The one-side shuffle of transform expressions came with SPARK-48012, in 4.0.0. So did the broadcast expansion of a keyed layout, since SPARK-49205 made `KeyGroupedPartitioning` an expression.

### Does this PR introduce _any_ user-facing change?

Yes, it fixes the wrong results, the crash and the planning failures above. Such a query now shuffles the side instead of pairing it or shuffling onto it.

That costs shuffles in three shapes whose plan was already correct:
- A widening cast counts as part of the argument's shape. When `t1.id` is `BIGINT` and `p.b` is `INT`, the first join reports `bucket(4, cast(b as bigint))`. A later join `p.b = t3.x` with an `INT` bucketed `t3` then shuffles both sides. Before, a connector whose bucket function has one canonical name for both types paired the two sides as they stood. This happens with every conf at its default, through a broadcast join. Two sides that both carry the cast still pair.
- Under `allowCompatibleTransforms`, an identity-partitioned side is no longer reduced onto such a transform, so it is shuffled.
- With `spark.sql.sources.v2.bucketing.shuffle.enabled`, a join on two keys whose sides pair only through a transform of an expression gets one more shuffle. Take two broadcast joins, `t1 JOIN /*+ BROADCAST(p) */ p ON t1.id = p.b + 1` and the same over `t2` and `q`, joined on `l.id = r.c AND l.b = r.b`. Each layout covers only one join key, so the storage-partitioned join declines under `requireAllClusterKeysForCoPartition`. Master then keeps both sides, paired on `bucket(4, b + 1)`. This PR cannot shuffle onto that layout, so it shuffles one side onto `bucket(4, id)`. Both plans co-partition on one key only, which `requireAllClusterKeysForCoPartition` is meant to prevent. That is a separate, pre-existing question, filed as SPARK-59971.

One plan gets better. A collection with a bare expression member, e.g. the `b + 1` of a side shuffled onto an identity-partitioned table, used to be refused as a layout as a whole. Now its other members can serve, so `ident i JOIN plain p ON i.id = p.b + 1 JOIN xy q ON i.id = q.x AND p.b = q.y` takes 2 shuffles instead of 3.

This PR keeps the fix small, since it has to reach the maintenance branches. SPARK-59900 is the follow-up for the rest: shuffling another side onto such a layout, and reducing an identity side onto it.

### How was this patch tested?

New tests. All but two of them fail on master:
- `ShuffleSpecSuite`, "SPARK-59887: a spec over a transform of an expression pairs only with the same shape". For an expression and a struct field argument, it covers the same shape, another shape and a column, both ways and with itself, the reduce in both directions, the two sides of one reduce, two sides reduced through different pairings, an identity side, `canCreatePartitioning`, and a collection with a usable sibling.
- `ShuffleSpecSuite`, "SPARK-59901: an expression over two columns maps to no cluster key".
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side shuffled onto a transform of an expression is not paired as it is". It runs the query above against a `bucket(4)` third table, and with `allowCompatibleTransforms` against a `bucket(8)` one that the join reduces, and asserts two shuffles.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side is not shuffled onto a transform of an expression". It runs both join orders and asserts that `xy` is shuffled onto `bucket(4, x)` and never onto a bucket of `y`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a broadcast join's layout over a join key is not paired as it is". It runs with every conf at its default and checks that the first join is a broadcast join.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a layout over a transform of a struct field is not paired nor shuffled onto". Master returns 0 of 8 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps what each member's keys are computed from". It covers a bucketed and an identity-partitioned first table, a bucketed and an unpartitioned third table, and the struct form. Master returns 0 of 7 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: an identity side is not reduced onto a transform of an expression". It runs with `p.b + 1` over a `Long.MaxValue` key and with `-p.b` over a `Long.MinValue` key. Master fails with the `ARITHMETIC_OVERFLOW` for both.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps a member over two columns". It runs the `GROUP BY` above. Master returns 4 groups instead of 2.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns is not paired through one of them". Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns does not cover its columns". It joins two broadcast joins that report `[b + c, d]` on `(b, c, d)` under `allowKeysSubsetOfPartitionKeys`, and asserts that the join is not paired on `d` alone. Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: two layouts over the same shape of a join key pair as two columns do". It runs the cast case above with every conf at its default. It also runs two one-side shuffles onto `bucket(4, b + 1)` and onto `bucket(4)` or `bucket(8)` of `q.b + 1`, where `allowCompatibleTransforms` reduces the `bucket(8)` one. It asserts that the join adds no shuffle. It passes on master, which pairs and reduces these as well.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: the same shape reduced together pairs only through the same pairing". Two legs each shuffle onto `bucket(12, b + 1)`. The first leg reduces with `bucket(8)`, so its keys are buckets of 4. The second leg reduces with `bucket(8)` too, with `bucket(18)`, or not at all, so its keys are buckets of 4, 6 or 12. It asserts 2 shuffles for the first and 4 for the other two. Master returns 0 of 23 rows for the third.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a member that cannot serve does not rank its collection". It pins the ranking: `q` keeps its layout with 3 partitions, as on master, where the collection is refused as a whole. Without the member filter, the collection would rank first with 5 partitions, and `q` would be shuffled onto 2.

The SPARK-59905 test from #59189 now also runs its `p.b + p.c` query with AQE on. Master fails that run with the `AssertionError`.

The three expression tests each run with the key `p.b + 1` and with `-p.b`. Master returns 0 of 8 rows for `p.b + 1` and 4 of 8 for `-p.b`. The `-p.b` key is the shape that also fails on 4.2, 4.1 and 4.0. On those branches `b + 1` already works, since its literal leaf blocks the pairing.

Each change was removed on its own:
- without the guard in `isExpressionCompatible`, the one-side shuffle and broadcast pairing tests, the struct field test and the reduce test return wrong rows, and the identity tests fail on the assertion in `rebuiltOver`;
- without the shape pairing, the same-shape and same-pairing tests fail on the shuffle count;
- with a reduce of such a pair refused, the same-shape test fails on the shuffle count;
- with every pair of reduced keys over an expression refused, the same-pairing test fails on the shuffle count;
- with the reduced marks ignored, or only checked on both sides, the same-pairing test returns 7 of 23 rows;
- without the `canCreatePartitioning` clause, the one-side shuffle pairing, `xy`, struct field, reduce and same-pairing tests fail, all on the assertion in `rebuiltOver`;
- with the `forall` back in `ShuffleSpecCollection.canCreatePartitioning`, the `xy` tests fail on the shuffle count;
- without the member filter in `pickCoPartitionTarget`, 15 tests fail in `KeyGroupedPartitioningSuite` and `EnsureRequirementsSuite`. The filter now also does what the collection-level check before it did, so 8 of them are older tests. Of this PR's tests, 7 fail, 6 of them on the assertion in `rebuiltOver`. A member over the same shape pairs with itself, so only the filter keeps it from being the layout;
- without the `GroupPartitionsExec` change, the reduce test returns 0 of 7 rows, and the two-column `GROUP BY` returns 4 groups;
- without the `requireAllClusterKeysForCoPartition` change, the coverage test pairs on `d` alone;
- with the `keyPositions` assertion back, both SPARK-59901 tests in `KeyGroupedPartitioningSuite` fail.

The first `ShuffleSpecSuite` test also fails for each change to `isExpressionCompatible` above.

The `xy` tests pass with the guard removed, since the `canCreatePartitioning` clause and the member filter keep `xy` off that layout.

Also ran `TransformExpressionSuite`, `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `WriteDistributionAndOrderingSuite`, `PlannerSuite`, `ProjectedOrderingAndPartitioningSuite`, `GroupPartitionsExecSuite`, `EnsureRequirementsSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite`, 1955 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 #59165 from peter-toth/SPARK-59887-spj-expression-key-identity.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
peter-toth added a commit that referenced this pull request Oct 6, 2026
…-partitioned join pairs a transform of a join-key expression with one of a column

### Differences from #59165

This is the branch-4.3 backport of #59165. The description below is the original one. On this branch:
- The cherry-pick conflicts in `partitioning.scala`, `GroupPartitionsExec`, `EnsureRequirements` and the suite's imports, since the code around the change differs on branch-4.3. The main code is ported by hand. The tests are master's, plus the ones named below.
- Branch-4.3 has no `pickCoPartitionTarget`. `EnsureRequirements` picks the layout through `bestSpecOpt` and `bestMemberOpt` instead. Only a member that can serve as the layout matches a child there, and `bestMemberOpt` picks only such a member.
- Branch-4.3 ranks a collection by its first member. Now that a collection can serve when any member can, that first member may be one that cannot. So a collection ranks by its first member that can serve, which is the branch's ranking whenever the first member can.
- Two new `EnsureRequirementsSuite` tests cover these two changes, since no master test reaches them on this branch. "SPARK-59887: a collection ranks by its first member that can serve as the layout" lists a member that cannot serve first. "SPARK-59887: a child matched only through a member that cannot serve is shuffled" pairs a child with the layout's other member only. Each fails with its change removed.
- Matching only through members that can serve costs a shuffle in one more shape on this branch. Two legs that line up only through members of the same shape now take a shuffle, e.g. `[a, bucket(4, -y)]` against `[m, bucket(4, -n)]`, joined on `l.a = r.m AND l.y = r.n`. With `allowKeysSubsetOfPartitionKeys`, the shuffle conf on and `requireAllClusterKeysForCoPartition` off, branch-4.3 took no shuffle, and this change shuffles one side onto 2 partitions. The rows are right either way.
  - Master pairs such legs in `checkKeyGroupCompatible`, which compares every member of both sides. It takes no shuffle there, before and after this change.
  - This branch's `createKeyedShuffleSpec` offers only the first member that satisfies the distribution, here `[a, bucket(4, c)]`. So only the shuffle path paired the legs, through a member that cannot serve as the layout. This change keeps the shuffle path from that, which is also where master's third cost above comes from.
  - A widening cast instead of the negation behaves the same. So "Two sides that both carry the cast still pair" holds here only where the storage-partitioned join path pairs them.
- Branch-4.3's `createKeyedShuffleSpec` checks `spark.sql.requireAllClusterKeysForCoPartition` by matching the layout's columns to the join keys one by one, where master checks coverage. So the coverage change becomes: a layout with an expression over two columns does not meet the requirement. With the rest of this change alone, the rows are right, but:
  - a join on `(b, c, d)` over `[b + c, d]` is paired on `d` alone;
  - a join on `(b, c)` over `[b + c]` is paired on no key at all, in one partition.
- That needs a sibling layout that satisfies the join's distribution, so that the child is not shuffled first. The coverage test, "SPARK-59901: a join key over two columns does not cover its columns", has none on branch-4.3, since each leg's projection drops the identity side's columns. So it also runs with those columns kept.
- `KeyedPartitioning.createShuffleSpec` gets master's refusal to project onto no position. It came with SPARK-59289, which is not on branch-4.3, and master's SPARK-59901 fix relies on it.
  - A join keyed on the expression itself, e.g. `l.b + l.c = r.k`, can satisfy the distribution under `spark.sql.requireAllClusterKeysForDistribution`, while no position maps to a cluster key.
  - With `allowKeysSubsetOfPartitionKeys`, the projection would then run the join in one partition. That needs `spark.sql.sources.v2.bucketing.shuffle.enabled`, or another side over the same expression with `requireAllClusterKeysForCoPartition` off. Any such key does it, e.g. `b + 1`, `-b`, or the cast of an `INT` key joined to a `BIGINT` one.
  - A `b + c` key failed there with the `AssertionError` before this change, and the others already ran in one partition. They now get the shuffle, which can mean 2 shuffles where branch-4.3 had none.
  - The new test "SPARK-59901: a join keyed on a join-key expression does not run in one partition" covers `b + 1` and `b + c`.
- `TransformExpressionSuite` does not exist on branch-4.3, so its hunk is left out. So are the doc hunks for code that branch-4.3 does not have.
- The bugs reproduce on branch-4.3 as on master, `b + 1` included, through both the one-side shuffle and the broadcast join. Measured without this change:
  - the `b + 1` query above returns 0 of 7 rows;
  - the `-p.b` shapes return 4 of 8 rows;
  - the struct field pairing returns 0 of 8 rows, and shuffling onto it fails with the `ClassCastException`.
- New tests: all but three fail on branch-4.3 without this change. The three are the two that also pass on master, and the `EnsureRequirementsSuite` ranking test, since branch-4.3 refuses that collection as a whole. The SPARK-59905 test's `p.b + p.c` run with AQE on passes without this change as well. The first query of "SPARK-59901: a join key over two columns is not paired through one of them" passes without this change as well, since AQE does not fail a single such join on this branch, so "a single join is enough" does not hold here. Its second query fails with the `AssertionError`.
- The main changes were removed one at a time: the shape gate, the `canCreatePartitioning` clause, the `exists` in `ShuffleSpecCollection.canCreatePartitioning`, the member filter, the ranking change, the matching change, the coverage change, the refusal to project onto no position, the `GroupPartitionsExec` change and the `keyPositions` assertion. Each fails the same tests as on master, except:
  - without the member filter, only the two `xy` tests fail, on the assertion in `rebuiltOver`;
  - with the `keyPositions` assertion back, the `ShuffleSpecSuite` SPARK-59901 test, "SPARK-59901: a join key over two columns is not paired through one of them" and "SPARK-59901: a join keyed on a join-key expression does not run in one partition" fail. The coverage test passes, since the coverage change keeps `[b + c, d]` away from `keyPositions`;
  - removing the ranking change, the matching change or the refusal to project onto no position fails only its new test. The master ranking test's collection lists its serving member first, so the branch's own ranking already gets that one right.
- Ran `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `ProjectedOrderingAndPartitioningSuite`, `PlannerSuite`, `WriteDistributionAndOrderingSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite` on branch-4.3, 1858 tests in all, plus `dev/lint-scala`.

### What changes were proposed in this pull request?

A storage-partitioned join now pairs a partition transform whose argument is not a bare column, i.e. is an expression or a struct field, only with a transform over the same argument. It also stops dropping such an argument where it used to:

- `KeyedShuffleSpec.isExpressionCompatible` pairs a transform with an argument that is not an `Attribute` only with a transform over the same argument shape, i.e. the argument with its one column replaced by a placeholder. So `bucket(4, b + 1)` pairs with `bucket(4, c + 1)` on `b = c`, but not with `bucket(4, x)`. Over the same shape the two compare as two transforms of columns do. Under `allowCompatibleTransforms`, `bucket(8, b + 1)` reduces onto `bucket(4, c + 1)`, since a reduce maps the partition keys and evaluates no argument. When an earlier join reduced the keys, two such transforms pair only when they were reduced through the same pairing. An identity side is never reduced onto such a transform. `TransformExpression` keeps its comparisons. It loses `withReference`, whose last callers now rebuild the transform with `copy`.
- `KeyedShuffleSpec.canCreatePartitioning` asks the same question, so such a layout is never the one other children are shuffled onto. `createPartitioning` replaces a transform's argument with the other child's cluster key, which drops whatever surrounds the key. It and the identity arm of `reducersBothWays` now rebuild a transform through one helper that asserts the argument is a column.
- `ShuffleSpecCollection.canCreatePartitioning` is now true when any member can, and `EnsureRequirements.pickCoPartitionTarget` offers and ranks only the members that can. So one member that cannot serve no longer rules out a usable sibling. That also covers a bare expression such as `b + 1`, which master already refused.
- `GroupPartitionsExec` rebuilds a reduced expression over each member's own argument when it reports the members of a collection. It used to re-target the column only.
- `KeyedShuffleSpec.keyPositions` maps an expression with more than one reference, e.g. `bucket(4, b + c)`, to no position instead of failing an assertion (SPARK-59901).
- `EnsureRequirements.createKeyedShuffleSpecs` counts only an expression over one column towards `spark.sql.requireAllClusterKeysForCoPartition`. An expression over two columns maps to no position, so a projection would drop it together with its columns.

### Why are the changes needed?

With `spark.sql.sources.v2.bucketing.shuffle.enabled` on, this query returns 0 rows instead of 7. `t1(id)` and `t3(x)` are partitioned by `bucket(4, ...)`, `plain(b)` is not partitioned, and each holds the values 0 to 7:

```sql
SELECT t1.id, p.b, t3.x FROM t1
JOIN plain p ON t1.id = p.b + 1
JOIN t3 ON p.b = t3.x
```

1. The first join shuffles `plain` onto `t1`'s partitioning. `KeyedShuffleSpec.createPartitioning` builds that partitioning from `plain`'s join key, so it is `bucket(4, b + 1)`. This join is correct.
2. The second join clusters `plain` on the bare `b`. `KeyedShuffleSpec.keyPositions` maps a partition expression to a cluster key through its reference, so `bucket(4, b + 1)` counts as a function of `b`, the counterpart of `x`.
3. `isSameFunction` compares only the function name and the bucket count. So `bucket(4, b + 1)` is the same as `t3`'s `bucket(4, x)`, and the join pairs the partitions as they stand. A row with `b = x` sits in bucket `(b + 1) % 4` on one side and in `x % 4` on the other.

The comparison ignores which column an argument is, and relies on `keyPositions` to pair the columns up. That is sound only when the argument is the column itself. `bucket(4, b + 1)` is a function of `b`, but not the same function of it as `bucket(4, x)` is of `x`. A scan never reports a transform of an expression or of a struct field, but two planner paths build one. A one-side shuffle does, under `spark.sql.sources.v2.bucketing.shuffle.enabled`. An inner broadcast hash join does too, with every conf at its default. `BroadcastHashJoinExec.expandOutputPartitioning` also reports the streamed side's layout over the build side's join key. The same problem shows up in more shapes, all measured on master:

- **A broadcast join.** With every conf at its default, `SELECT /*+ BROADCAST(p) */ ... FROM t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x` returns 0 of 7 rows the same way, through the expanded `bucket(4, b + 1)` layout.
- **The reducer path.** Under `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`, a `bucket(8, x)` third table loses the rows the same way.
- **A struct field.** A join key `p.s.a` gives `bucket(4, s.a)`, which `keyPositions` pairs through `s`. Joining two such sides on `s` pairs `bucket(4, s.a)` with `bucket(4, s.b)` and returns 0 of 8 rows. Shuffling another side onto such a layout builds `bucket(4, s)`, which fails on executors with a `ClassCastException`.
- **Shuffling onto the layout.** In `plain p JOIN t1 ON p.b + 1 = t1.id JOIN xy q ON t1.id = q.x AND p.b = q.y`, the second join shuffles `xy` onto the first join's `bucket(4, b + 1)` layout. That builds `bucket(4, y)` for `xy`, drops the `+ 1`, and returns 0 of 8 rows.
- **After a reduce.** When a later join reduces the first join's output from `bucket(4)` onto `bucket(2)`, `GroupPartitionsExec` reports the `bucket(4, b + 1)` member as `bucket(2, b)`. A join on `b` then pairs it with a `bucket(2, x)` scan, or shuffles another side onto it, and returns 0 of 7 rows. The same happens to the bare `b + 1` member of a side shuffled onto an identity-partitioned table. The struct form fails with the `ClassCastException`.
- **An identity side reduced onto it.** Under `allowCompatibleTransforms`, a later join can reduce an identity-partitioned side onto `bucket(4, b + 1)`. That evaluates `bucket(4, id + 1)` on the side's partition keys while planning. The keys come out right, but under ANSI a key of `Long.MaxValue` fails the query with `ARITHMETIC_OVERFLOW`, although the query never computes `id + 1`.
- **A GROUP BY after a reduce.** With `allowCompatibleTransforms`, `SELECT /*+ BROADCAST(p) */ p.b, max(p.c) FROM ident i JOIN plain2 p ON i.id = p.b + p.c JOIN bucket2 t2 ON i.id = t2.id GROUP BY p.b` returns 4 groups instead of 2. The reduce reports the `p.b + p.c` member as `bucket(2, p.b)`, which serves the `GROUP BY` as it stands.
- **A key over two columns.** A join key such as `p.b + p.c` gives `bucket(4, b + c)`, or the bare `b + c` over an identity-partitioned side. Any spec over it fails planning with an `AssertionError` in `keyPositions` (SPARK-59901). AQE validates the plan in `OptimizeSkewedJoin`, so a single join is enough.

An `ORDER BY` over such a partition expression had a related problem in another code path. That was SPARK-59905, fixed in #59189.

The one-side shuffle of transform expressions came with SPARK-48012, in 4.0.0. So did the broadcast expansion of a keyed layout, since SPARK-49205 made `KeyGroupedPartitioning` an expression.

### Does this PR introduce _any_ user-facing change?

Yes, it fixes the wrong results, the crash and the planning failures above. Such a query now shuffles the side instead of pairing it or shuffling onto it.

That costs shuffles in three shapes whose plan was already correct:
- A widening cast counts as part of the argument's shape. When `t1.id` is `BIGINT` and `p.b` is `INT`, the first join reports `bucket(4, cast(b as bigint))`. A later join `p.b = t3.x` with an `INT` bucketed `t3` then shuffles both sides. Before, a connector whose bucket function has one canonical name for both types paired the two sides as they stood. This happens with every conf at its default, through a broadcast join. Two sides that both carry the cast still pair.
- Under `allowCompatibleTransforms`, an identity-partitioned side is no longer reduced onto such a transform, so it is shuffled.
- With `spark.sql.sources.v2.bucketing.shuffle.enabled`, a join on two keys whose sides pair only through a transform of an expression gets one more shuffle. Take two broadcast joins, `t1 JOIN /*+ BROADCAST(p) */ p ON t1.id = p.b + 1` and the same over `t2` and `q`, joined on `l.id = r.c AND l.b = r.b`. Each layout covers only one join key, so the storage-partitioned join declines under `requireAllClusterKeysForCoPartition`. Master then keeps both sides, paired on `bucket(4, b + 1)`. This PR cannot shuffle onto that layout, so it shuffles one side onto `bucket(4, id)`. Both plans co-partition on one key only, which `requireAllClusterKeysForCoPartition` is meant to prevent. That is a separate, pre-existing question, filed as SPARK-59971.

One plan gets better. A collection with a bare expression member, e.g. the `b + 1` of a side shuffled onto an identity-partitioned table, used to be refused as a layout as a whole. Now its other members can serve, so `ident i JOIN plain p ON i.id = p.b + 1 JOIN xy q ON i.id = q.x AND p.b = q.y` takes 2 shuffles instead of 3.

This PR keeps the fix small, since it has to reach the maintenance branches. SPARK-59900 is the follow-up for the rest: shuffling another side onto such a layout, and reducing an identity side onto it.

### How was this patch tested?

New tests. All but two of them fail on master:
- `ShuffleSpecSuite`, "SPARK-59887: a spec over a transform of an expression pairs only with the same shape". For an expression and a struct field argument, it covers the same shape, another shape and a column, both ways and with itself, the reduce in both directions, the two sides of one reduce, two sides reduced through different pairings, an identity side, `canCreatePartitioning`, and a collection with a usable sibling.
- `ShuffleSpecSuite`, "SPARK-59901: an expression over two columns maps to no cluster key".
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side shuffled onto a transform of an expression is not paired as it is". It runs the query above against a `bucket(4)` third table, and with `allowCompatibleTransforms` against a `bucket(8)` one that the join reduces, and asserts two shuffles.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side is not shuffled onto a transform of an expression". It runs both join orders and asserts that `xy` is shuffled onto `bucket(4, x)` and never onto a bucket of `y`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a broadcast join's layout over a join key is not paired as it is". It runs with every conf at its default and checks that the first join is a broadcast join.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a layout over a transform of a struct field is not paired nor shuffled onto". Master returns 0 of 8 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps what each member's keys are computed from". It covers a bucketed and an identity-partitioned first table, a bucketed and an unpartitioned third table, and the struct form. Master returns 0 of 7 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: an identity side is not reduced onto a transform of an expression". It runs with `p.b + 1` over a `Long.MaxValue` key and with `-p.b` over a `Long.MinValue` key. Master fails with the `ARITHMETIC_OVERFLOW` for both.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps a member over two columns". It runs the `GROUP BY` above. Master returns 4 groups instead of 2.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns is not paired through one of them". Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns does not cover its columns". It joins two broadcast joins that report `[b + c, d]` on `(b, c, d)` under `allowKeysSubsetOfPartitionKeys`, and asserts that the join is not paired on `d` alone. Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: two layouts over the same shape of a join key pair as two columns do". It runs the cast case above with every conf at its default. It also runs two one-side shuffles onto `bucket(4, b + 1)` and onto `bucket(4)` or `bucket(8)` of `q.b + 1`, where `allowCompatibleTransforms` reduces the `bucket(8)` one. It asserts that the join adds no shuffle. It passes on master, which pairs and reduces these as well.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: the same shape reduced together pairs only through the same pairing". Two legs each shuffle onto `bucket(12, b + 1)`. The first leg reduces with `bucket(8)`, so its keys are buckets of 4. The second leg reduces with `bucket(8)` too, with `bucket(18)`, or not at all, so its keys are buckets of 4, 6 or 12. It asserts 2 shuffles for the first and 4 for the other two. Master returns 0 of 23 rows for the third.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a member that cannot serve does not rank its collection". It pins the ranking: `q` keeps its layout with 3 partitions, as on master, where the collection is refused as a whole. Without the member filter, the collection would rank first with 5 partitions, and `q` would be shuffled onto 2.

The SPARK-59905 test from #59189 now also runs its `p.b + p.c` query with AQE on. Master fails that run with the `AssertionError`.

The three expression tests each run with the key `p.b + 1` and with `-p.b`. Master returns 0 of 8 rows for `p.b + 1` and 4 of 8 for `-p.b`. The `-p.b` key is the shape that also fails on 4.2, 4.1 and 4.0. On those branches `b + 1` already works, since its literal leaf blocks the pairing.

Each change was removed on its own:
- without the guard in `isExpressionCompatible`, the one-side shuffle and broadcast pairing tests, the struct field test and the reduce test return wrong rows, and the identity tests fail on the assertion in `rebuiltOver`;
- without the shape pairing, the same-shape and same-pairing tests fail on the shuffle count;
- with a reduce of such a pair refused, the same-shape test fails on the shuffle count;
- with every pair of reduced keys over an expression refused, the same-pairing test fails on the shuffle count;
- with the reduced marks ignored, or only checked on both sides, the same-pairing test returns 7 of 23 rows;
- without the `canCreatePartitioning` clause, the one-side shuffle pairing, `xy`, struct field, reduce and same-pairing tests fail, all on the assertion in `rebuiltOver`;
- with the `forall` back in `ShuffleSpecCollection.canCreatePartitioning`, the `xy` tests fail on the shuffle count;
- without the member filter in `pickCoPartitionTarget`, 15 tests fail in `KeyGroupedPartitioningSuite` and `EnsureRequirementsSuite`. The filter now also does what the collection-level check before it did, so 8 of them are older tests. Of this PR's tests, 7 fail, 6 of them on the assertion in `rebuiltOver`. A member over the same shape pairs with itself, so only the filter keeps it from being the layout;
- without the `GroupPartitionsExec` change, the reduce test returns 0 of 7 rows, and the two-column `GROUP BY` returns 4 groups;
- without the `requireAllClusterKeysForCoPartition` change, the coverage test pairs on `d` alone;
- with the `keyPositions` assertion back, both SPARK-59901 tests in `KeyGroupedPartitioningSuite` fail.

The first `ShuffleSpecSuite` test also fails for each change to `isExpressionCompatible` above.

The `xy` tests pass with the guard removed, since the `canCreatePartitioning` clause and the member filter keep `xy` off that layout.

Also ran `TransformExpressionSuite`, `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `WriteDistributionAndOrderingSuite`, `PlannerSuite`, `ProjectedOrderingAndPartitioningSuite`, `GroupPartitionsExecSuite`, `EnsureRequirementsSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite`, 1955 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 #59237 from peter-toth/SPARK-59887-spj-expression-key-identity-4.3.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
peter-toth added a commit that referenced this pull request Oct 6, 2026
…-partitioned join pairs a transform of a join-key expression with one of a column

### Differences from #59165

This is the branch-4.2 backport of #59165. The description below is the original one. The cherry-pick does not apply, since the storage-partitioned join code on this branch predates master's rework of it. So the fix is ported onto this branch's own code:
- `KeyedShuffleSpec.keyPositions` maps a partition expression by its leaves here, not by its references. So a key with a literal, e.g. `b + 1`, already maps to no cluster key and is not paired wrongly on this branch. A single-leaf key such as `-b` or `s.a` is. The assertion in `keyPositions` fails here too for an expression with two leaves, `b + 1` included (SPARK-59901). It is replaced by "only an expression with a single leaf maps to a position".
- This branch picks the layout to shuffle onto in `EnsureRequirements` through `candidateSpecs`, `bestSpecOpt` and `bestMemberOpt`, not `pickCoPartitionTarget`. With `ShuffleSpecCollection.canCreatePartitioning` an `exists`, a collection whose first member cannot serve becomes a candidate here. So three places look at its serving members only:
  - the ranking counts a collection by its first serving member, which is its first member whenever that one serves, as before;
  - a child counts as matched only through a serving member;
  - the layout is picked from the serving members.
- Two new tests in `EnsureRequirementsSuite` cover the ranking and the matching, since master's tests do not reach them on this branch: "SPARK-59887: a collection ranks by its first member that can serve as the layout" and "SPARK-59887: a child matched only through a member that cannot serve is shuffled". Each fails with its change removed.
- Matching only through members that can serve costs a shuffle in one more shape on this branch. Two legs that line up only through members of the same shape now take a shuffle, e.g. `[a, bucket(4, -y)]` against `[m, bucket(4, -n)]`, joined on `l.a = r.m AND l.y = r.n`. With `allowJoinKeysSubsetOfPartitionKeys`, the shuffle conf on and `requireAllClusterKeysForCoPartition` off, branch-4.2 took no shuffle, and this change shuffles one side onto 2 partitions. The rows are right either way.
  - Master pairs such legs in `checkKeyGroupCompatible`, which compares every member of both sides. It takes no shuffle there, before and after this change.
  - This branch's `createKeyedShuffleSpec` offers only the first member that satisfies the distribution, here `[a, bucket(4, c)]`. So only the shuffle path paired the legs, through a member that cannot serve as the layout. This change keeps the shuffle path from that, which is also where master's third cost above comes from.
  - A widening cast instead of the negation behaves the same. So "Two sides that both carry the cast still pair" holds here only where the storage-partitioned join path pairs them.
- This branch's `createKeyedShuffleSpec` checks `spark.sql.requireAllClusterKeysForCoPartition` by matching the leaves of the layout to the join keys one by one, where master checks coverage. So the coverage change becomes: a layout with an expression with several leaves does not meet the requirement. With the rest of this change alone, the rows are right, but a join on `(b, c, d)` over `[b + c, d]` is paired on `d` alone.
- That needs a sibling layout that satisfies the join's distribution, so that the child is not shuffled first. The coverage test, "SPARK-59901: a join key over two columns does not cover its columns", has none on this branch, since each leg's projection drops the identity side's columns. So it also runs with those columns kept.
- `KeyedPartitioning.createShuffleSpec` gets master's refusal to project onto no position. It came with SPARK-59289, which is not on this branch, and master's SPARK-59901 fix relies on it.
  - A join keyed on the expression itself, e.g. `l.b + l.c = r.k`, can satisfy the distribution under `spark.sql.requireAllClusterKeysForDistribution`, while no position maps to a cluster key.
  - With `allowJoinKeysSubsetOfPartitionKeys`, the projection would then run the join in one partition. That needs `spark.sql.sources.v2.bucketing.shuffle.enabled`, or another side over the same expression with `requireAllClusterKeysForCoPartition` off. A `-b` key already did that on this branch.
  - A `b + 1` or `b + c` key failed there with the `AssertionError` before this change, even with the shuffle conf off. Now all three get the shuffle, which can mean 2 shuffles where branch-4.2 had none.
  - The new test "SPARK-59901: a join keyed on a join-key expression does not run in one partition" covers `-b`, `b + 1` and `b + c`.
- `TransformExpressionSuite` does not exist on this branch, so its hunk is dropped. So are the doc changes for code this branch does not have.
- Three tests are tailored, since master's version cannot pass here. On this branch a key with a literal maps to no cluster key, so it is shuffled rather than paired:
  - the first `ShuffleSpecSuite` test uses `-a`, `-b` and `abs(b)` for master's `a + 1`, `b + 1` and `b + 2`;
  - "two layouts over the same shape of a join key pair as two columns do" and "the same shape reduced together pairs only through the same pairing" use `-p.b` over negative values for master's `p.b + 1`.
- Measured on branch-4.2 without the change, with the shuffle conf on: the `-p.b` key returns 4 of 8 rows through a one-side shuffle and through the `bucket(8)` reducer path. Through a broadcast join it returns 4 of 8 rows with every conf at its default. The struct field pairing returns 0 of 8 rows, and shuffling onto it fails with a `ClassCastException`.
- Two parts of the original description need a `-p.b` key on this branch, since a member with a literal does not satisfy the distribution here:
  - the plan that gets better: `ident i JOIN plain p ON i.id = -p.b JOIN xy q ON i.id = q.x AND p.b = q.y` takes 2 shuffles instead of 3. With `p.b + 1` it already took 2;
  - the third cost: with `-p.b` branch-4.2 took no shuffle, and this change takes 1. With `p.b + 1` one side was already shuffled.
- A single join keyed on `p.b + p.c` does not fail on this branch, AQE included, so "a single join is enough" does not hold here. So the first query of "SPARK-59901: a join key over two columns is not paired through one of them" passes without the change, and its second query fails with the `AssertionError`.
- These included tests pass on branch-4.2 without the change:
  - the `p.b + 1` variants of the shuffled-onto, `xy`, broadcast and identity tests;
  - "two layouts over the same shape of a join key pair as two columns do";
  - "a member that cannot serve does not rank its collection", since its `p.y + 1` member does not satisfy the distribution here;
  - the SPARK-59905 test, which now runs `p.b + p.c` with AQE on as on master;
  - the new ranking test, since a collection with a member that cannot serve is no candidate there.
- Each change was removed on its own:
  - without the shape gate, the `-p.b` shuffled-onto, broadcast and identity tests and the struct field test fail;
  - without the `canCreatePartitioning` clause, the `-p.b` shuffled-onto and `xy` tests, the struct field, reduce and same-pairing tests fail on the assertion in `rebuiltOver`. The new matching test and the first `ShuffleSpecSuite` test fail too;
  - without the member filter, the `-p.b` `xy` test and both new `EnsureRequirementsSuite` tests fail;
  - with the ranking by the first member back, the new ranking test fails;
  - with the matching through any member back, the new matching test fails;
  - without the `GroupPartitionsExec` change, the reduce test and the two-column `GROUP BY` fail;
  - with the `keyPositions` assertion back, the `ShuffleSpecSuite` SPARK-59901 test, "SPARK-59901: a join key over two columns is not paired through one of them" and "SPARK-59901: a join keyed on a join-key expression does not run in one partition" fail. The coverage test passes, since the coverage change keeps `[b + c, d]` away from `keyPositions`;
  - with the `forall` back in `ShuffleSpecCollection.canCreatePartitioning`, the `-p.b` `xy` test fails on the shuffle count, and so do the new matching test and the first `ShuffleSpecSuite` test;
  - removing the coverage change or the refusal to project onto no position fails only its new test or run.
- Ran `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `EnsureRequirementsSuite`, `GroupPartitionsExecSuite`, `ProjectedOrderingAndPartitioningSuite`, `PlannerSuite`, `WriteDistributionAndOrderingSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite` on branch-4.2, 1403 tests in all, plus `dev/lint-scala`.

### What changes were proposed in this pull request?

A storage-partitioned join now pairs a partition transform whose argument is not a bare column, i.e. is an expression or a struct field, only with a transform over the same argument. It also stops dropping such an argument where it used to:

- `KeyedShuffleSpec.isExpressionCompatible` pairs a transform with an argument that is not an `Attribute` only with a transform over the same argument shape, i.e. the argument with its one column replaced by a placeholder. So `bucket(4, b + 1)` pairs with `bucket(4, c + 1)` on `b = c`, but not with `bucket(4, x)`. Over the same shape the two compare as two transforms of columns do. Under `allowCompatibleTransforms`, `bucket(8, b + 1)` reduces onto `bucket(4, c + 1)`, since a reduce maps the partition keys and evaluates no argument. When an earlier join reduced the keys, two such transforms pair only when they were reduced through the same pairing. An identity side is never reduced onto such a transform. `TransformExpression` keeps its comparisons. It loses `withReference`, whose last callers now rebuild the transform with `copy`.
- `KeyedShuffleSpec.canCreatePartitioning` asks the same question, so such a layout is never the one other children are shuffled onto. `createPartitioning` replaces a transform's argument with the other child's cluster key, which drops whatever surrounds the key. It and the identity arm of `reducersBothWays` now rebuild a transform through one helper that asserts the argument is a column.
- `ShuffleSpecCollection.canCreatePartitioning` is now true when any member can, and `EnsureRequirements.pickCoPartitionTarget` offers and ranks only the members that can. So one member that cannot serve no longer rules out a usable sibling. That also covers a bare expression such as `b + 1`, which master already refused.
- `GroupPartitionsExec` rebuilds a reduced expression over each member's own argument when it reports the members of a collection. It used to re-target the column only.
- `KeyedShuffleSpec.keyPositions` maps an expression with more than one reference, e.g. `bucket(4, b + c)`, to no position instead of failing an assertion (SPARK-59901).
- `EnsureRequirements.createKeyedShuffleSpecs` counts only an expression over one column towards `spark.sql.requireAllClusterKeysForCoPartition`. An expression over two columns maps to no position, so a projection would drop it together with its columns.

### Why are the changes needed?

With `spark.sql.sources.v2.bucketing.shuffle.enabled` on, this query returns 0 rows instead of 7. `t1(id)` and `t3(x)` are partitioned by `bucket(4, ...)`, `plain(b)` is not partitioned, and each holds the values 0 to 7:

```sql
SELECT t1.id, p.b, t3.x FROM t1
JOIN plain p ON t1.id = p.b + 1
JOIN t3 ON p.b = t3.x
```

1. The first join shuffles `plain` onto `t1`'s partitioning. `KeyedShuffleSpec.createPartitioning` builds that partitioning from `plain`'s join key, so it is `bucket(4, b + 1)`. This join is correct.
2. The second join clusters `plain` on the bare `b`. `KeyedShuffleSpec.keyPositions` maps a partition expression to a cluster key through its reference, so `bucket(4, b + 1)` counts as a function of `b`, the counterpart of `x`.
3. `isSameFunction` compares only the function name and the bucket count. So `bucket(4, b + 1)` is the same as `t3`'s `bucket(4, x)`, and the join pairs the partitions as they stand. A row with `b = x` sits in bucket `(b + 1) % 4` on one side and in `x % 4` on the other.

The comparison ignores which column an argument is, and relies on `keyPositions` to pair the columns up. That is sound only when the argument is the column itself. `bucket(4, b + 1)` is a function of `b`, but not the same function of it as `bucket(4, x)` is of `x`. A scan never reports a transform of an expression or of a struct field, but two planner paths build one. A one-side shuffle does, under `spark.sql.sources.v2.bucketing.shuffle.enabled`. An inner broadcast hash join does too, with every conf at its default. `BroadcastHashJoinExec.expandOutputPartitioning` also reports the streamed side's layout over the build side's join key. The same problem shows up in more shapes, all measured on master:

- **A broadcast join.** With every conf at its default, `SELECT /*+ BROADCAST(p) */ ... FROM t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x` returns 0 of 7 rows the same way, through the expanded `bucket(4, b + 1)` layout.
- **The reducer path.** Under `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`, a `bucket(8, x)` third table loses the rows the same way.
- **A struct field.** A join key `p.s.a` gives `bucket(4, s.a)`, which `keyPositions` pairs through `s`. Joining two such sides on `s` pairs `bucket(4, s.a)` with `bucket(4, s.b)` and returns 0 of 8 rows. Shuffling another side onto such a layout builds `bucket(4, s)`, which fails on executors with a `ClassCastException`.
- **Shuffling onto the layout.** In `plain p JOIN t1 ON p.b + 1 = t1.id JOIN xy q ON t1.id = q.x AND p.b = q.y`, the second join shuffles `xy` onto the first join's `bucket(4, b + 1)` layout. That builds `bucket(4, y)` for `xy`, drops the `+ 1`, and returns 0 of 8 rows.
- **After a reduce.** When a later join reduces the first join's output from `bucket(4)` onto `bucket(2)`, `GroupPartitionsExec` reports the `bucket(4, b + 1)` member as `bucket(2, b)`. A join on `b` then pairs it with a `bucket(2, x)` scan, or shuffles another side onto it, and returns 0 of 7 rows. The same happens to the bare `b + 1` member of a side shuffled onto an identity-partitioned table. The struct form fails with the `ClassCastException`.
- **An identity side reduced onto it.** Under `allowCompatibleTransforms`, a later join can reduce an identity-partitioned side onto `bucket(4, b + 1)`. That evaluates `bucket(4, id + 1)` on the side's partition keys while planning. The keys come out right, but under ANSI a key of `Long.MaxValue` fails the query with `ARITHMETIC_OVERFLOW`, although the query never computes `id + 1`.
- **A GROUP BY after a reduce.** With `allowCompatibleTransforms`, `SELECT /*+ BROADCAST(p) */ p.b, max(p.c) FROM ident i JOIN plain2 p ON i.id = p.b + p.c JOIN bucket2 t2 ON i.id = t2.id GROUP BY p.b` returns 4 groups instead of 2. The reduce reports the `p.b + p.c` member as `bucket(2, p.b)`, which serves the `GROUP BY` as it stands.
- **A key over two columns.** A join key such as `p.b + p.c` gives `bucket(4, b + c)`, or the bare `b + c` over an identity-partitioned side. Any spec over it fails planning with an `AssertionError` in `keyPositions` (SPARK-59901). AQE validates the plan in `OptimizeSkewedJoin`, so a single join is enough.

An `ORDER BY` over such a partition expression had a related problem in another code path. That was SPARK-59905, fixed in #59189.

The one-side shuffle of transform expressions came with SPARK-48012, in 4.0.0. So did the broadcast expansion of a keyed layout, since SPARK-49205 made `KeyGroupedPartitioning` an expression.

### Does this PR introduce _any_ user-facing change?

Yes, it fixes the wrong results, the crash and the planning failures above. Such a query now shuffles the side instead of pairing it or shuffling onto it.

That costs shuffles in three shapes whose plan was already correct:
- A widening cast counts as part of the argument's shape. When `t1.id` is `BIGINT` and `p.b` is `INT`, the first join reports `bucket(4, cast(b as bigint))`. A later join `p.b = t3.x` with an `INT` bucketed `t3` then shuffles both sides. Before, a connector whose bucket function has one canonical name for both types paired the two sides as they stood. This happens with every conf at its default, through a broadcast join. Two sides that both carry the cast still pair.
- Under `allowCompatibleTransforms`, an identity-partitioned side is no longer reduced onto such a transform, so it is shuffled.
- With `spark.sql.sources.v2.bucketing.shuffle.enabled`, a join on two keys whose sides pair only through a transform of an expression gets one more shuffle. Take two broadcast joins, `t1 JOIN /*+ BROADCAST(p) */ p ON t1.id = p.b + 1` and the same over `t2` and `q`, joined on `l.id = r.c AND l.b = r.b`. Each layout covers only one join key, so the storage-partitioned join declines under `requireAllClusterKeysForCoPartition`. Master then keeps both sides, paired on `bucket(4, b + 1)`. This PR cannot shuffle onto that layout, so it shuffles one side onto `bucket(4, id)`. Both plans co-partition on one key only, which `requireAllClusterKeysForCoPartition` is meant to prevent. That is a separate, pre-existing question, filed as SPARK-59971.

One plan gets better. A collection with a bare expression member, e.g. the `b + 1` of a side shuffled onto an identity-partitioned table, used to be refused as a layout as a whole. Now its other members can serve, so `ident i JOIN plain p ON i.id = p.b + 1 JOIN xy q ON i.id = q.x AND p.b = q.y` takes 2 shuffles instead of 3.

This PR keeps the fix small, since it has to reach the maintenance branches. SPARK-59900 is the follow-up for the rest: shuffling another side onto such a layout, and reducing an identity side onto it.

### How was this patch tested?

New tests. All but two of them fail on master:
- `ShuffleSpecSuite`, "SPARK-59887: a spec over a transform of an expression pairs only with the same shape". For an expression and a struct field argument, it covers the same shape, another shape and a column, both ways and with itself, the reduce in both directions, the two sides of one reduce, two sides reduced through different pairings, an identity side, `canCreatePartitioning`, and a collection with a usable sibling.
- `ShuffleSpecSuite`, "SPARK-59901: an expression over two columns maps to no cluster key".
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side shuffled onto a transform of an expression is not paired as it is". It runs the query above against a `bucket(4)` third table, and with `allowCompatibleTransforms` against a `bucket(8)` one that the join reduces, and asserts two shuffles.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a side is not shuffled onto a transform of an expression". It runs both join orders and asserts that `xy` is shuffled onto `bucket(4, x)` and never onto a bucket of `y`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a broadcast join's layout over a join key is not paired as it is". It runs with every conf at its default and checks that the first join is a broadcast join.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a layout over a transform of a struct field is not paired nor shuffled onto". Master returns 0 of 8 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps what each member's keys are computed from". It covers a bucketed and an identity-partitioned first table, a bucketed and an unpartitioned third table, and the struct form. Master returns 0 of 7 rows.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: an identity side is not reduced onto a transform of an expression". It runs with `p.b + 1` over a `Long.MaxValue` key and with `-p.b` over a `Long.MinValue` key. Master fails with the `ARITHMETIC_OVERFLOW` for both.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps a member over two columns". It runs the `GROUP BY` above. Master returns 4 groups instead of 2.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns is not paired through one of them". Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns does not cover its columns". It joins two broadcast joins that report `[b + c, d]` on `(b, c, d)` under `allowKeysSubsetOfPartitionKeys`, and asserts that the join is not paired on `d` alone. Master fails with the `AssertionError`.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: two layouts over the same shape of a join key pair as two columns do". It runs the cast case above with every conf at its default. It also runs two one-side shuffles onto `bucket(4, b + 1)` and onto `bucket(4)` or `bucket(8)` of `q.b + 1`, where `allowCompatibleTransforms` reduces the `bucket(8)` one. It asserts that the join adds no shuffle. It passes on master, which pairs and reduces these as well.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: the same shape reduced together pairs only through the same pairing". Two legs each shuffle onto `bucket(12, b + 1)`. The first leg reduces with `bucket(8)`, so its keys are buckets of 4. The second leg reduces with `bucket(8)` too, with `bucket(18)`, or not at all, so its keys are buckets of 4, 6 or 12. It asserts 2 shuffles for the first and 4 for the other two. Master returns 0 of 23 rows for the third.
- `KeyGroupedPartitioningSuite`, "SPARK-59887: a member that cannot serve does not rank its collection". It pins the ranking: `q` keeps its layout with 3 partitions, as on master, where the collection is refused as a whole. Without the member filter, the collection would rank first with 5 partitions, and `q` would be shuffled onto 2.

The SPARK-59905 test from #59189 now also runs its `p.b + p.c` query with AQE on. Master fails that run with the `AssertionError`.

The three expression tests each run with the key `p.b + 1` and with `-p.b`. Master returns 0 of 8 rows for `p.b + 1` and 4 of 8 for `-p.b`. The `-p.b` key is the shape that also fails on 4.2, 4.1 and 4.0. On those branches `b + 1` already works, since its literal leaf blocks the pairing.

Each change was removed on its own:
- without the guard in `isExpressionCompatible`, the one-side shuffle and broadcast pairing tests, the struct field test and the reduce test return wrong rows, and the identity tests fail on the assertion in `rebuiltOver`;
- without the shape pairing, the same-shape and same-pairing tests fail on the shuffle count;
- with a reduce of such a pair refused, the same-shape test fails on the shuffle count;
- with every pair of reduced keys over an expression refused, the same-pairing test fails on the shuffle count;
- with the reduced marks ignored, or only checked on both sides, the same-pairing test returns 7 of 23 rows;
- without the `canCreatePartitioning` clause, the one-side shuffle pairing, `xy`, struct field, reduce and same-pairing tests fail, all on the assertion in `rebuiltOver`;
- with the `forall` back in `ShuffleSpecCollection.canCreatePartitioning`, the `xy` tests fail on the shuffle count;
- without the member filter in `pickCoPartitionTarget`, 15 tests fail in `KeyGroupedPartitioningSuite` and `EnsureRequirementsSuite`. The filter now also does what the collection-level check before it did, so 8 of them are older tests. Of this PR's tests, 7 fail, 6 of them on the assertion in `rebuiltOver`. A member over the same shape pairs with itself, so only the filter keeps it from being the layout;
- without the `GroupPartitionsExec` change, the reduce test returns 0 of 7 rows, and the two-column `GROUP BY` returns 4 groups;
- without the `requireAllClusterKeysForCoPartition` change, the coverage test pairs on `d` alone;
- with the `keyPositions` assertion back, both SPARK-59901 tests in `KeyGroupedPartitioningSuite` fail.

The first `ShuffleSpecSuite` test also fails for each change to `isExpressionCompatible` above.

The `xy` tests pass with the guard removed, since the `canCreatePartitioning` clause and the member filter keep `xy` off that layout.

Also ran `TransformExpressionSuite`, `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `WriteDistributionAndOrderingSuite`, `PlannerSuite`, `ProjectedOrderingAndPartitioningSuite`, `GroupPartitionsExecSuite`, `EnsureRequirementsSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite`, 1955 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 #59238 from peter-toth/SPARK-59887-spj-expression-key-identity-4.2.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
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.

3 participants