Skip to content

[SPARK-59887][SPARK-59901][SQL][4.3] Fix wrong results when a storage-partitioned join pairs a transform of a join-key expression with one of a column - #59237

Closed
peter-toth wants to merge 1 commit into
apache:branch-4.3from
peter-toth:SPARK-59887-spj-expression-key-identity-4.3
Closed

peter-toth wants to merge 1 commit into
apache:branch-4.3from
peter-toth:SPARK-59887-spj-expression-key-identity-4.3

Conversation

@peter-toth

Copy link
Copy Markdown
Contributor

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:

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)

…-partitioned join pairs a transform of a join-key expression with one of a column

### Differences from apache#59165

This is the branch-4.3 backport of apache#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 the same as on master.
- 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.
- The `spark.sql.requireAllClusterKeysForCoPartition` change is left out. Here a member over two columns never satisfies a clustered distribution under `allowKeysSubsetOfPartitionKeys`, so its child is shuffled before the co-partitioning check. Its test, "SPARK-59901: a join key over two columns does not cover its columns", is kept, and it passes on this branch without this change.
- `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 coverage test above. Both new `ShuffleSpecSuite` tests fail too. The SPARK-59905 test's `p.b + p.c` run with AQE on passes without this change as well.
- The main changes were removed one at a time: the shape gate, the `canCreatePartitioning` clause, the `forall`, the member filter, the ranking change, the matching change, 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, only "SPARK-59901: a join key over two columns is not paired through one of them" fails among the `KeyGroupedPartitioningSuite` tests;
  - removing the ranking change or the matching change fails only its new `EnsureRequirementsSuite` 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, 1857 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 apache#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 two 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 the first join's output. Before, a connector whose bucket function has one canonical name for both types paired it as it 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.

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 apache#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)
@uros-b

uros-b commented Oct 6, 2026

Copy link
Copy Markdown
Member

Thank you @peter-toth and @dongjoon-hyun!

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 peter-toth closed this Oct 6, 2026
@peter-toth

Copy link
Copy Markdown
Contributor Author

Merge Summary:

Posted by merge_spark_pr.py

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