Skip to content

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

Open
peter-toth wants to merge 4 commits into
apache:masterfrom
peter-toth:SPARK-59887-spj-expression-key-identity
Open

peter-toth wants to merge 4 commits into
apache:masterfrom
peter-toth:SPARK-59887-spj-expression-key-identity

Conversation

@peter-toth

@peter-toth peter-toth commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

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

  • KeyedShuffleSpec.isExpressionCompatible refuses any pair in which a side is a transform with an argument that is not an Attribute. That covers a pair as it stands, the reduce path, a pair an earlier join reduced together, and an identity side reduced onto such a transform. TransformExpression itself is unchanged.
  • 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.
  • ShuffleSpecCollection.canCreatePartitioning is now true when any member can, and EnsureRequirements.pickCoPartitionTarget offers only the members that can. So one such member no longer rules out a usable sibling.
  • 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).

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, 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 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 a partition expression that a one-side shuffle built has a related problem in another code path. That is SPARK-59905.

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 a few shapes whose plan was already correct:

  • When two sides carry the same expression shape and no other layout, e.g. bucket(4, b + 1) and bucket(4, c + 1) joined on b = c, both sides are shuffled.
  • An identity-partitioned side is no longer reduced onto such a transform, so it is shuffled.
  • A widening cast counts as an expression. When t1.id is BIGINT and p.b is INT, the first join builds 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.
  • A scan that reports a bucket of a struct field, e.g. bucket(4, s.a), no longer pairs on a join over the whole struct. The same check keeps it from pairing bucket(4, s.a) with bucket(4, t.b). The in-memory test catalog reports no such partitioning, so no test covers this.

This PR is deliberately the minimal fix. It fixes wrong results that go back to 4.0.0, so it has to reach the maintenance branches, and a refusal is the smallest change that is safe to backport. Once it is in, I plan a follow-up, SPARK-59900, to give these shuffles back. It describes a transform's argument relative to its column, so that b + 1 and c + 1 compare equal while b + 1 and x do not.

How was this patch tested?

New tests, each of which fails on master:

  • ShuffleSpecSuite, "SPARK-59887: a spec over a transform of an expression pairs with nothing". For an expression and a struct field argument, it covers isCompatibleWith both ways and with itself, the reduce in both directions, a pair reduced together, 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". Master fails with the ARITHMETIC_OVERFLOW.
  • KeyGroupedPartitioningSuite, "SPARK-59901: a join key over two columns is not paired through one of them". Master fails 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 isExpressionCompatible guard, the three pairing tests, the struct field test and the reduce test return wrong rows, and the identity test fails with the ARITHMETIC_OVERFLOW;
  • without the canCreatePartitioning clause, the struct field and reduce tests fail with the ClassCastException, and the pairing tests fail on the shuffle count;
  • without the GroupPartitionsExec change, the reduce test returns 0 of 7 rows;
  • with the keyPositions assertion back, both SPARK-59901 tests fail;
  • with the forall back, the xy tests fail on the shuffle count.

Removing only the member filter in pickCoPartitionTarget changes no test. The guard already keeps such a member from pairing with its own side, so the filter is defensive.

The xy tests pass with either the guard or the canCreatePartitioning clause removed. Each keeps xy off that layout on its own. The guard does it because the layout then does not pair with its own side, which pickCoPartitionTarget requires of the member it picks.

Also ran TransformExpressionSuite, ShuffleSpecSuite, DistributionSuite, the KeyGroupedPartitioning* suites, WriteDistributionAndOrderingSuite, PlannerSuite, ProjectedOrderingAndPartitioningSuite, GroupPartitionsExecSuite, EnsureRequirementsSuite, BucketedReadWithoutHiveSupportSuite, the plan stability suites, the *JoinSuite suites and AdaptiveQueryExecSuite, 1945 tests in all, plus dev/lint-scala.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Code

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

`TransformExpression.isSameFunction` and `TransformExpression.reducers` now say no when either transform has an argument that is not a bare column, an `Attribute`. `isCompatible` goes through `reducers`, so the one check covers the reducer path too. `KeyedShuffleSpec.canCreatePartitioning` asks the same question, so such a layout is never the one other children are shuffled onto.

With `spark.sql.sources.v2.bucketing.shuffle.enabled` on, `KeyedShuffleSpec.createPartitioning` can shuffle a side onto a transform of its join key, such as `bucket(4, b + 1)` or `bucket(4, s.a)`. `keyPositions` maps such an expression to a cluster key through its reference, and `isSameFunction` ignores the arguments. So a later join can pair it with a different function of the same key, or shuffle another side onto it with the expression around the key dropped. Either way the join loses rows.

@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 left 13 inline comments on the latest commit. I ran items 1-4 and 6 as scratch tests in KeyGroupedPartitioningSuite, on both this PR and its base commit. Here is a summary, ordered by priority.

Correctness

  1. TransformExpression L104: the new guard is bypassed once a later join reduces the keys. GroupPartitionsExec re-targets a bucket(4, b + 1) member as bucket(2, b), so the query returns 0 of 7 rows, and the struct form fails with a ClassCastException.
  2. KeyedShuffleSpec L2205: a two-column join-key expression such as p.b + p.c fails planning with an AssertionError in keyPositions, even in a single join. The new clause is evaluated after it.
  3. KeyedShuffleSpec L1948: the same reference-based reading makes ORDER BY 100 - p.b return rows in reverse order.

Performance

  1. KeyedShuffleSpec L2212: ShuffleSpecCollection.canCreatePartitioning is a forall, so two shapes take 4 shuffles instead of 2, more than the PR description says.
  2. TransformExpression L75: an implicit widening Cast on the join key is refused too.

Tests

  1. KeyGroupedPartitioningSuite L1576: the xy test fails on the base commit only because plain is the left side of the first join.
  2. TransformExpressionSuite L150: reducers and isCompatible are not tested for a struct-field argument.

Design and cleanup

  1. TransformExpression L133: one guard in KeyedShuffleSpec.isExpressionCompatible would cover every production path and keep isSameFunction reflexive.
  2. KeyedShuffleSpec L2208: the underlying issue is that createPartitioning replaces the whole argument. This is input for SPARK-59900.
  3. TransformExpression L117: the isCompatible and reducers Scaladocs, and the reflexivity note of ShuffleSpec.isCompatibleWith, still describe the old contract.
  4. TransformExpression L74: the Literal case is unreachable, and its rationale doesn't hold for either consumer.
  5. KeyGroupedPartitioningSuite L1609: the two struct-field tests duplicate their setup.
  6. TransformExpressionSuite L45: DivisibleBucketFunction duplicates ShuffleSpecSuite.FakeBucket.

Items 1-3 also happen on the base commit, so they predate this PR. Still, items 1 and 2 matter most here: item 1 reaches the wrong result that this PR fixes through a path the new guard does not cover, and item 2 is a planning failure in the same createPartitioning shape.

*/
def isSameFunction(other: TransformExpression): Boolean = functionId == other.functionId
def isSameFunction(other: TransformExpression): Boolean =
argumentsAreAttributes && other.argumentsAreAttributes && functionId == other.functionId

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 new guard is bypassed once a later join reduces the keys, and the wrong result comes back. GroupPartitionsExec rewrites every member of the reduced child's partitioning with the reduced expression re-targeted at the member's first reference:

case (expr, Some(KeyReducer(_, reduced))) =>
// `reduced` came from the one member `checkKeyGroupCompatible` paired
// this side on, which need not be the member being rewritten. The keys are
// reduced once, from the shared key rows, so `reduced` describes them
// whichever member this is, and only the key attribute is re-targeted.
reduced.withReference(expr.references.head)

So a bucket(4, b + 1) member, which the one-side shuffle created, is reported as bucket(2, b) after a join reduces bucket(4) onto bucket(2). Its argument is a bare column, so isSameFunction pairs it and canCreatePartitioning accepts it, while its rows still sit in (b + 1) % 2. The following returns 0 of 7 rows on both this PR and its base commit:

Seq("bucket4" -> 4, "bucket2a" -> 2, "bucket2b" -> 2).foreach { case (name, n) =>
  createBucketedIdTable(name, n, numIds = 8)
}
createTable("plain", Array(Column.create("b", LongType)), Array.empty)
sql("INSERT INTO testcat.ns.plain VALUES " + (0 until 8).map(b => s"($b)").mkString(", "))
withSQLConf(
    SQLConf.V2_BUCKETING_SHUFFLE_ENABLED.key -> "true",
    SQLConf.V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS.key -> "true") {
  val df = sql(
    """SELECT b4.id, p.b, t2.id, t3.id FROM testcat.ns.bucket4 b4
      |JOIN testcat.ns.plain p ON b4.id = p.b + 1
      |JOIN testcat.ns.bucket2a t2 ON b4.id = t2.id
      |JOIN testcat.ns.bucket2b t3 ON p.b = t3.id""".stripMargin)
  checkAnswer(df, (0 until 7).map(b => Row(b + 1L, b.toLong, b + 1L, b.toLong)))
}

The second join reduces the left side's bucket(4, id) onto bucket(2), and the third join pairs the re-targeted bucket(2, b) with bucket(2, id) without any exchange. The same happens in three more shapes:

  • When the third table is unpartitioned, it is shuffled onto bucket(2, y), and the query returns 0 of 7 rows.
  • When the first table is identity-partitioned, the identity reduce turns the bare b + 1 member into bucket(2, b), and the query returns 0 of 7 rows.
  • In the struct form, bucket(4, s.a) becomes bucket(2, s), and another side shuffled onto it fails with the same ClassCastException as the struct-field tests of this PR.

Could GroupPartitionsExec re-target by the member's own argument instead? That means replacing the argument of reduced with the member's argument, i.e. the child of a transform, or the member itself for a bare expression. The member is then reported as bucket(2, b + 1), which describes its keys and which the new guard refuses.

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 fixed in 69f16ab the way you suggest. GroupPartitionsExec now rebuilds reduced over the member's own argument, the child of a transform or the member itself. So a bucket(4, b + 1) member becomes bucket(2, b + 1), which the guard refuses. reduced's own argument is always a bare column now, since the guard also covers the identity arm (thread).

"SPARK-59887: a reduce keeps what each member's keys are computed from" runs your four shapes. It returns 0 of 7 rows without the change.

@@ -2199,8 +2203,14 @@ case class KeyedShuffleSpec(
// cannot rewrite it. This also keeps the unprojected spec returned by `createShuffleSpec`
// for a marked narrowing projection from being chosen as the best spec.
keyPositions.forall(_.nonEmpty) &&

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.

A two-column join-key expression fails planning in keyPositions, before the new clause below is reached. createPartitioning builds the partition expression over the other child's join key, which can reference two columns, e.g. bucket(4, b + c), or the bare b + c for an identity-partitioned side. Any spec over it then fails the assertion at L1958. This single join fails on both this PR and its base commit:

createTable("ident", Array(Column.create("id", LongType)), Array(identity("id")))
sql("INSERT INTO testcat.ns.ident VALUES " + (0 until 8).map(i => s"($i)").mkString(", "))
createTable("plain2", Array(Column.create("b", LongType), Column.create("c", LongType)),
  Array.empty)
sql("INSERT INTO testcat.ns.plain2 VALUES (0, 1), (0, 2), (1, 1), (1, 2)")
withSQLConf(SQLConf.V2_BUCKETING_SHUFFLE_ENABLED.key -> "true") {
  sql("SELECT * FROM testcat.ns.ident i JOIN testcat.ns.plain2 p ON i.id = p.b + p.c")
    .collect()
}
[INTERNAL_ERROR] The Spark SQL phase planning failed with an internal error. ...
Cause: java.lang.AssertionError: assertion failed: Expected exactly one child from
(b#10L + c#11L), but found 2

It is reached through AQE's OptimizeSkewedJoin, which validates the plan with ValidateRequirements, so it needs nothing but spark.sql.sources.v2.bucketing.shuffle.enabled. With a bucketed side, bucket4 b4 JOIN plain2 p ON b4.id = p.b + p.c JOIN xy q ON p.b = q.x AND p.c = q.y fails the same way in EnsureRequirements.pickCoPartitionTarget, through this line.

I know this predates this PR, but it is the same createPartitioning shape. Moving the new clause ahead of keyPositions fixes only the second path. So createPartitioning, or EnsureRequirements when it picks the target, would need to refuse a cluster key with more than one reference. Alternatively, keyPositions could return an empty position set instead of asserting.

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.

Fixed in 69f16ab, as SPARK-59901, which I had filed for this and have now added to the title. keyPositions maps an expression with more than one reference to no position instead of asserting. canCreatePartitioning and areKeysCompatible both turn an empty position away, so a spec over b + c pairs with nothing and is no layout to shuffle onto. Both of your queries are in "SPARK-59901: a join key over two columns is not paired through one of them".

* "some reference is a cluster key" and "every reference is" are the same statement.
*
* This says which cluster key an expression is a function of, not which function it is:
* `bucket(4, b + 1)` maps to `b` just as `bucket(4, b)` does. Telling those two apart is

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 reference-based reading also breaks ORDER BY. When EnsureRequirements checks a keyed child against an OrderedDistribution, it binds the required ordering to expressions.flatMap(_.references):

// The single-column invariant in KeyedPartitioning.supportsExpressions guarantees
// one attribute per partition expression.
val attrs = satisfyingKeyedPartitioning.expressions.flatMap(_.references)
val keyRowOrdering = RowOrdering.create(o.ordering, attrs)

For a bare expression member that createPartitioning built over an identity side, e.g. 100 - b, the key rows already hold 100 - b. The ordering then computes 100 - (100 - b) and sorts the partitions the wrong way. The global sort is removed, and the rows come back in reverse order on both this PR and its base commit:

// `ident` is identity-partitioned on `id` and holds 94..100, and `plain(b)` holds 0..6.
withSQLConf(
    SQLConf.V2_BUCKETING_SHUFFLE_ENABLED.key -> "true",
    SQLConf.V2_BUCKETING_SORTING_ENABLED.key -> "true") {
  val df = sql(
    """SELECT p.b FROM testcat.ns.ident i JOIN testcat.ns.plain p ON i.id = 100 - p.b
      |ORDER BY 100 - p.b""".stripMargin)
  // Returns 0, 1, ..., 6 instead of 6, 5, ..., 0.
  assert(df.collect().map(_.getLong(0)).toSeq == (0 to 6).reverse.map(_.toLong))
}

This isn't from this PR either. But it is one more consumer that relies on the argument being the column itself, like the one this paragraph describes, so it may be worth covering here or in a separate JIRA.

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.

Confirmed. It is a separate consumer, so I filed SPARK-59905 for it rather than growing this PR. The same binding also fails planning for ORDER BY p.b + p.c over the two-column key, with an ArrayIndexOutOfBoundsException, so I added that shape to the ticket. The EnsureRequirements comment that claimed one attribute per partition expression for this binding now points there.

// which drops whatever surrounds the key, e.g. the `+ 1` of `bucket(4, b + 1)` or the `.a`
// of `bucket(4, s.a)`. Such a spec's own child would then be laid out differently from the
// children shuffled onto it.
case t: TransformExpression => t.argumentsAreAttributes

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.

One refused member disables the whole collection, so the cost is larger than the PR description says. ShuffleSpecCollection.canCreatePartitioning is a forall, so this clause also rules out the usable bucket(4, id) member of the same join output. The keyed layout is then lost for every operator above it too. I measured these with spark.sql.sources.v2.bucketing.shuffle.enabled on, and both return correct results:

  • The example in the PR description, t1 JOIN plain p ON t1.id = p.b + 1 JOIN xy q ON t1.id = q.x AND p.b = q.y, followed by GROUP BY t1.id, takes 4 shuffles instead of 2. xy and the first join's output are hash-shuffled on (x, y) and (id, b), and the aggregate then needs another exchange on id.
  • Two one-side-shuffled legs, (t1 JOIN p ON t1.id = p.b + 1) JOIN (t2 JOIN q ON t2.id = q.c + 1) ON t1.id = t2.id AND p.b = q.c, take 4 shuffles instead of 2. On the base commit, the outer join pairs the two legs as they stand.

Could the refusal be per member instead? For instance, ShuffleSpecCollection.canCreatePartitioning could be an exists, and EnsureRequirements.pickCoPartitionTarget could keep only the winner's members whose own canCreatePartitioning holds when it builds bestMembers. Then bucket(4, id) stays usable, and bucket(4, b + 1) is never used as a template. The change is as local as this clause, so it seems safe to backport too.

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, done in 69f16ab. ShuffleSpecCollection.canCreatePartitioning is now an exists, and pickCoPartitionTarget offers only the winner's members that can serve. PartitioningCollection.createShuffleSpec admits only members that could satisfy the distribution, so a collection without a keyed member hardly ever mixes the two answers. The planner, join, bucketed read, AQE and plan stability suites pass unchanged. Your two shapes and the xy one take 2 shuffles again, and the description no longer lists them as costs.

*/
private[sql] lazy val argumentsAreAttributes: Boolean = children.forall {
case _: Literal => true
case c => c.isInstanceOf[Attribute]

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.

This also refuses an implicit widening Cast on the join key. When t1.id is BIGINT and p.b is INT, type coercion casts only the INT side, so the first join shuffles p onto bucket(4, cast(b as bigint)). A second join p.b = t3.x, with t3 bucketed on an INT column x, was paired as it stood before this PR. For a connector whose bucket function has one canonical name for INT and BIGINT, that was a correct plan, and with this PR it costs one more shuffle of the first join's output. The PR description doesn't list this shape among the costs. Connectors with type-specific canonical names are not affected, since their functions never compared equal. If this matters, a widening cast of an attribute that Cast.canUpCast allows could count as the column.

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.

Right. I added it to the costs in the description. I'd rather not treat an up-cast as the column in a fix meant for the maintenance branches. That leans on each connector's canonical name meaning the same across input types.

e.isInstanceOf[AttributeReference] || e.isInstanceOf[TransformExpression]
partitioning.expressions.forall {
case _: AttributeReference => true
// `createPartitioning` replaces a transform's argument with the other child's cluster key,

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 underlying issue is that createPartitioning replaces the whole argument (L2225). This clause refuses a template because the rewrite drops whatever surrounds the key. Substituting only the reference that keyPositions mapped, as withReference and the IdentityReducer case already do, would be sound for a single-reference argument. bucket(4, b + 1) would become bucket(4, y + 1) for xy, and bucket(4, s.a) would become bucket(4, t.a), which also removes the ClassCastException. It is not a drop-in replacement for this clause, though:

  • To get the conceded shuffles back, isSameFunction would also have to compare the shapes and stay reflexive, since pickCoPartitionTarget keeps the winning child in place only when its member pairs with itself.
  • The shuffled side would evaluate y + 1 on rows where the query never computes it, and that can overflow under ANSI mode.

So I'm raising it as input for SPARK-59900 rather than as a change to this PR.

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, this is very useful for SPARK-59900, the ANSI point especially. Substituting the reference makes the shuffled side evaluate y + 1 on rows where the query never computes it.

}
}
def isCompatible(other: TransformExpression): Boolean =
isSameFunction(other) || reducers(other).isDefined || other.reducers(this).isDefined

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 Scaladocs still describe the old contract. The isCompatible Scaladoc right above says it is true if both functions are ReducibleFunctions and a Reducer exists, and the reducers Scaladoc says a reducer exists whenever one function is reducible on the other. The new test contradicts both: bucket8.reducers(bucket4) is defined, but bucket8OfExpression.isCompatible(bucket4) is false for the same functions. In addition, ShuffleSpec.isCompatibleWith still says "Spark assumes this to be reflexive, symmetric and transitive". It already had an exception in RangeShuffleSpec, and now a spec over bucket(4, b + 1) is incompatible with itself too. pickCoPartitionTarget avoids such a winner only because of the new canCreatePartitioning clause. Could you mention the argument condition in these docs?

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 69f16ab. With the guard out of TransformExpression, its docs hold again as they were. The ShuffleSpec.isCompatibleWith note now lists the specs that are not reflexive, and says what keeps the planner from picking one as a layout.

* transform rather than a column, so it does not count.
*/
private[sql] lazy val argumentsAreAttributes: Boolean = children.forall {
case _: Literal => 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.

The Literal case is unreachable, and its rationale doesn't hold for either consumer. No transform with a literal child reaches isSameFunction, reducers or canCreatePartitioning. A scan's transform has to pass KeyedPartitioning.supportsExpressions, which requires a single reference child, and createPartitioning keeps the arity and fills the argument with a join key. If one did arrive, e.g. truncate(a, 2), both consumers would get it wrong. functionId ignores the literal's value, so truncate(a, 2) would be the same function as truncate(b, 5). And createPartitioning would overwrite the literal and build truncate(x, x), which is the hazard the new comment in canCreatePartitioning describes. How about children.forall(_.isInstanceOf[Attribute]), without the sentence about literals? It is equivalent today, and it fails safe.

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 69f16ab. The check is children.forall(_.isInstanceOf[Attribute]), without the literal case.

}
}

test("SPARK-59887: a side is not shuffled onto a transform of a struct field") {

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: This test repeats the setup of the previous one (bucket4 with 8 ids and the same named_struct rows) and its l subquery. Per the ablation in the PR description, it fails only when the canCreatePartitioning clause is removed, and the previous test fails then too. If you want to keep this ClassCastException shape, it could be a second query in the previous test, SELECT l.id FROM <the same l> JOIN testcat.ns.qs r ON l.s = r.s, like createKeyTables() shares the setup of the expression tests. That saves the rs table and about 20 lines.

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.

Merged into one test in 69f16ab. The ClassCastException shape is its second query, over a plain ps.

}

/** Reduces a bucket count onto any smaller count that divides it. */
private object DivisibleBucketFunction

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: ShuffleSpecSuite.FakeBucket already reduces a bucket count onto a divisor with % otherNumBuckets, and its doc says it is local only because catalyst has no BucketFunction. Moving it next to FlipLowBitFunction in sql/catalyst/src/test/scala/org/apache/spark/sql/connector/catalog/functions/ would let both suites share it. Also, only the positive branch of this reducer runs, since the new guard in reducers refuses the other cases before the function is called.

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.

Gone in 69f16ab. The unit test moved to ShuffleSpecSuite, where it uses FakeBucket.

- Refuse a transform of an expression at the entry of KeyedShuffleSpec.isExpressionCompatible, which covers every pairing arm, the identity reduce included. TransformExpression is back to master.
- GroupPartitionsExec rebuilds a reduced expression over each member's own argument.
- keyPositions maps an expression with more than one reference to no position (SPARK-59901).
- ShuffleSpecCollection.canCreatePartitioning is true when any member can, and pickCoPartitionTarget offers only those members.
- Tests and docs.
@peter-toth peter-toth changed the title [SPARK-59887][SQL] Fix wrong results when a storage-partitioned join pairs a transform of a join-key expression with one of a column [SPARK-59887][SPARK-59901][SQL] Fix wrong results when a storage-partitioned join pairs a transform of a join-key expression with one of a column Oct 1, 2026
@peter-toth

Copy link
Copy Markdown
Contributor Author

Thanks a lot for the thorough review, @dongjoon-hyun. 69f16ab takes all of it except item 3, now SPARK-59905, and item 5, which the description lists as a cost. It also takes SPARK-59901 in, for item 2. I updated the title and the description.

@uros-b uros-b left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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

@peter-toth

Copy link
Copy Markdown
Contributor Author

@dongjoon-hyun, your broadcast join point on #59189 (comment) applies here too. BroadcastHashJoinExec.expandOutputPartitioning reports bucket(4, id) as bucket(4, p.b + 1) as well, so this PR's wrong result does not need spark.sql.sources.v2.bucketing.shuffle.enabled. 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 on master.

This PR already returns all 7 rows, since the guard does not depend on where the layout came from. 5375cc3 adds a test for it, and I updated the description.

@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 addressing round 1, @peter-toth, and for adding the broadcast join test in 5375cc3. I reviewed the latest commit, 5375cc3, and left 14 inline comments. This time I traced the scenarios through the code at 5375cc3 and at the base commit, but did not run them. Here is a summary, ordered by priority.

Performance

  1. KeyedShuffleSpec L2077: the listed costs also reach plans with default confs through the broadcast expansion. For example, two broadcast joins on an INT dimension key over BIGINT buckets, joined on that key, go from 0 to 2 fact-sized shuffles.
  2. EnsureRequirements L1128: a member over two columns passes requireAllClusterKeysForCoPartition and is then projected away, so a join can pair on a subset of its keys.

Tests

  1. GroupPartitionsExec L611: the rebuild also fixes a GROUP BY or window over a two-column member, where the base returns duplicate groups, but no test covers it.
  2. KeyGroupedPartitioningSuite L1701: the identity-reduce test uses only p.b + 1, which branch-4.2 never pairs, so it would pass there without the fix.
  3. EnsureRequirements L379: the member filter also changes the ranking in a reachable shape, so it is not only defensive, and no test pins it.

Docs and comments

  1. KeyedShuffleSpec L1948: the keyPositions Scaladoc names only createPartitioning as a producer, and overstates when createShuffleSpec projects.
  2. ShuffleSpec L1590: following R1-10, the list of non-reflexive specs misses specs whose keyPositions entry is empty, e.g. years(ts) clustered on years(ts).
  3. IdentityReducer L1897: this comment and step 3.3 of isCompatibleWith do not reflect the argument condition.
  4. EnsureRequirements L142: two copies of the single-column comment remain, at L534 and L541.
  5. GroupPartitionsExec L615: the Scaladoc of TransformExpression.withReference and a test comment still describe the use this line removes.

API and design

  1. KeyedShuffleSpec L1965: mapping keyPositions by the argument would make the three shape checks agree by construction. This may be input for SPARK-59900.

Cleanup

  1. EnsureRequirements L376: the collection guard repeats the member filter, and the .max at L391 relies on that.
  2. ShuffleSpecSuite L778: the new bucket helper reverses the argument order used elsewhere in the suite, and the s = b example does not type-check.
  3. KeyGroupedPartitioningSuite L1688: plain is created inline three times, with an unexplained 0 until 7.
  4. The PR description: the Generated-by line has no tool version, which the template asks for. I left no inline comment for this one.

Items 1, 3 and 4 matter most for the backport. Item 1 changes plans with default confs on the maintenance branches, and items 3 and 4 are paths this PR fixes that no test would catch, on master and on branch-4.2 respectively.

right: Expression,
allowReduce: Boolean): Boolean = {
if (TransformExpression.hasReducedKeys(left) || TransformExpression.hasReducedKeys(right)) {
val overExpression = Seq(left, right).exists {

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 costs listed in the description now also reach plans that use only default confs, through the broadcast expansion. The description says an inner broadcast hash join builds such a layout with every conf at its default, but the costs still read as if they need spark.sql.sources.v2.bucketing.shuffle.enabled. For example, take orders and returns bucketed by bucket(16, cust_id) on a BIGINT column, and an INT customers.id broadcast into both:

SELECT a.id FROM
  (SELECT c1.id FROM orders o JOIN customers c1 ON o.cust_id = c1.id) a
JOIN
  (SELECT c2.id FROM returns r JOIN customers c2 ON r.cust_id = c2.id) b
ON a.id = b.id

Each broadcast join also reports bucket(16, cast(cN.id as bigint)). The base pairs the two as they stand, with no shuffle, and that is correct since both sides have the same shape. With this guard, both fact-sized sides are hash-shuffled, since keyed shuffles are off by default. The t.id = d.k + 1 shape behaves the same. It is the first cost in the list, but on the maintenance branches it changes plans for users who set nothing. Could the description say that these costs apply with default confs too?

def allClusterKeysCovered: Boolean =
// The single-column invariant in KeyedPartitioning.supportsExpressions guarantees one
// attribute per partition expression.
// Every column a partition expression references counts as covered. A scan reports one

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.

A member over two columns passes this check, and KeyedPartitioning.createShuffleSpec then projects it away. This counts every column that p.b + p.c references as covered, but keyPositions now maps p.b + p.c to no position (partitioning.scala:1964-1968). So under spark.sql.sources.v2.bucketing.allowKeysSubsetOfPartitionKeys.enabled, createShuffleSpec drops that position (partitioning.scala:1015). For instance, take two sides that each report [p.b + p.c, p.d], e.g. the output of a broadcast join ON i.x = p.b + p.c AND i.d = p.d over a table partitioned by identity(x), identity(d). A join of the two on (b, c, d) passes this check with {b, c, d}, but both specs are projected to [p.d], and the storage-partitioned join pairs them on d alone. The rows are right, but the join runs on as many partitions as there are distinct values of d, which is the skew spark.sql.requireAllClusterKeysForCoPartition is there to prevent. On the base the same plan failed the keyPositions assertion, so it is a newly reachable plan rather than a regression. Could this count only the columns of expressions with a single reference, or check the projected spec instead?

// admits no other. The member's own argument takes its place, since that need
// not be one. A side shuffled onto this layout reports `bucket(8, b + 1)` next
// to `bucket(8, id)`, or `b + 1` next to an identity `id`.
val argument = expr match {

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.

This also fixes a GROUP BY or window over a member with two columns, but no test covers that. The base re-targeted a p.b + p.c member at its first reference, bucket(2, p.b), which satisfies ClusteredDistribution([p.b]), so a GROUP BY p.b or PARTITION BY p.b above it skipped its shuffle. With spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled on, plain2 holding (0, 1), (0, 2), (1, 1), (1, 2), and the ident and bucket2 tables of the new tests:

sql("""SELECT /*+ BROADCAST(p) */ p.b, max(p.c) FROM testcat.ns.ident i
      |JOIN testcat.ns.plain2 p ON i.id = p.b + p.c
      |JOIN testcat.ns.bucket2 t2 ON i.id = t2.id
      |GROUP BY p.b""".stripMargin)

The base returns (0, 2), (1, 1), (0, 1), (1, 2), since the rows of b = 0 sit in two buckets, and this PR returns (0, 2), (1, 2). The aggregate has to keep p.c: with count(*), column pruning drops the member before the second join. The reduce test (KeyGroupedPartitioningSuite.scala:1647) covers single-column members under a later join only, and the SPARK-59901 test (KeyGroupedPartitioningSuite.scala:1707) has no reduce. Could one of them add this shape, or its row_number() OVER (PARTITION BY p.b ORDER BY p.c) variant?

SQLConf.V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS.key -> "true") {
val df = sql(
"""SELECT k.id, p.b, i.id FROM testcat.ns.bucket4 k
|JOIN testcat.ns.plain p ON k.id = p.b + 1

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.

This test would pass on branch-4.2 without the fix. It is the only end-to-end test of the identity arm, and it runs with p.b + 1 only. As the description says, branch-4.2 never pairs p.b + 1, because its createKeyedShuffleSpec compares the leaves of the partition expressions with the clustering and sees the literal. But branch-4.2 has the identity arm too, IdentityReducer(t.withReference(a)), so -p.b does pair with ident there, and the planner evaluates bucket(4, -id) on the keys of ident. Under ANSI that fails with ARITHMETIC_OVERFLOW for a key of Long.MinValue, although the query never negates it. The two expression tests above run both keys for the same reason. Could this test run -p.b as well, with Long.MinValue in ident?

(!shouldConsiderMinParallelism ||
child.outputPartitioning.numPartitions >= conf.defaultNumShufflePartitions) =>
child -> spec
child -> spec.flatten.filter(_.canCreatePartitioning)

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.

This filter is not only defensive: it can change which child the layout comes from. The description says that removing it changes no test, since the guard already keeps a refused member from pairing with its own side. That holds for the pairing, but the ranking at L391 now takes the max over the serving members only, and with spark.sql.sources.v2.bucketing.allowKeysSubsetOfPartitionKeys.enabled the members of one collection can project to different counts. For example, with that conf and spark.sql.sources.v2.bucketing.shuffle.enabled on, take t1 partitioned by identity(a), identity(c) with keys (1, 1), (1, 2), (1, 3), (2, 4), (2, 5), and SELECT /*+ BROADCAST(p) */ * FROM t1 JOIN p ON t1.c = p.y + 1 JOIN q ON t1.a = q.m AND p.y = q.n, where q reports 3 partitions without an exchange. The broadcast join reports [a, c] and [a, y + 1]. The first projects to [a], 2 partitions, and can serve. The second keeps 5 and cannot. With the filter, q wins with 3, as on the base. Without it, t1 ranks at 5 and wins, and q is shuffled onto 2 partitions. Both plans are correct, so only the parallelism differs. Could the description say so, and could a test in EnsureRequirementsSuite pin the ranking?

case t: TransformExpression => t.children
case e => Seq(e)
}
reduced.copy(children = argument)

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: this replaces the use that TransformExpression.withReference was added for in SPARK-59045, but its Scaladoc (TransformExpression.scala:127-131) still relies on the single leaf of KeyedPartitioning.supportsExpressions and warns about field paths. TransformExpressionSuite.scala:115-116 also still says the marker survives the rewrites that "GroupPartitionsExec apply". The only callers left are the identity arm (partitioning.scala:2207 and partitioning.scala:2211), which sees a transform over a single bare column. Could the identity arm use t.copy(children = Seq(a)) like this line, so that withReference and its test can go? Otherwise, could the two comments be updated?

assert(refs.size == 1, s"Expected exactly one child from $e, but found ${refs.size}")
distKeyToPos.getOrElse(refs.head.canonicalized, mutable.BitSet.empty)
if (refs.size == 1) {
distKeyToPos.getOrElse(refs.head.canonicalized, mutable.BitSet.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.

Minor: how about mapping by the argument here instead of by the reference? createPartitioning replaces the argument of a transform, and GroupPartitionsExec now rebuilds over it (GroupPartitionsExec.scala:611-615), but this maps an expression by its single reference. If it looked up t.children.head for a single-child transform, and the expression itself otherwise, every shape the new tests refuse would map to no position: bucket(4, a + 1) under [a], bucket(4, s.x) under [s], a + b under [a, b], and the reduced bucket(2, p.b + 1) under [p.b]. The overlap test in areKeysCompatible, canCreatePartitioning and createPartitioning would then agree by construction, without argumentsAreColumns or the refs.size == 1 branch. Today three shape checks have to be kept in step by hand, and they already differ: refs.size == 1 here, the leaf match after overExpression in isExpressionCompatible, and AttributeReference or argumentsAreColumns in canCreatePartitioning. Unlike the SPARK-59900 idea, this evaluates nothing new, so the ANSI concern does not arise. The identity arm at L2207 and L2211 would have to switch to t.copy(children = Seq(a)) at the same time. Otherwise, with spark.sql.requireAllClusterKeysForCoPartition off, a join keyed on p.b + 1 would give bucket(4, p.b + 1) a position, and withReference would build bucket(4, id + 1). I understand this may fit SPARK-59900 better, so please take it as input.

// that can serve as the layout. Any other member would build a partitioning its own child is
// not laid out on.
val candidateSpecs = children.zip(specs).collect {
case (child, Some(spec)) if spec.canCreatePartitioning &&

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: spec.canCreatePartitioning here is the same as the member list below being non-empty, since a collection now answers exists over its members. So one decision has two sources, and the .max at L391 depends on that identity in partitioning.scala. If ShuffleSpecCollection.canCreatePartitioning changes, the .max throws on an empty list, which the SPARK-59256 test comment warns about. Could this drop the guard and keep the children whose members are non-empty instead? The comments at L385-L386 ("the best any member offers") and L405 ("over all of its members") could then say serving members.

KeyedShuffleSpec(
KeyedPartitioning(Seq(expression), Seq(InternalRow(0L), InternalRow(1L))),
ClusteredDistribution(expression.references.toSeq))
def bucket(argument: Expression, numBuckets: Int): TransformExpression =

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: this reverses the argument order of bucket(numBuckets, expr) at L676, of the connector Expressions.bucket(numBuckets, ...), and of the comments of this test, e.g. L782 says bucket(8, argument) while L784 calls bucket(argument, 8). Also, the s = b example at L766-L767 joins a struct with a long, which the analyzer rejects. The shape that can happen is a join of two structs, s = t, which pairs bucket(4, s.a) with bucket(4, t.b). Could the helper take the bucket count first, and the comment use s = t?

test("SPARK-59887: an identity side is not reduced onto a transform of an expression") {
createBucketedIdTable("bucket4", 4, numIds = 8)
createTable("plain", Array(Column.create("b", LongType)), Array.empty)
sql("INSERT INTO testcat.ns.plain VALUES " + (0 until 7).map(b => s"($b)").mkString(", "))

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: plain is created inline three times in the new tests (L1543-L1544, L1650-L1651 and here), and only this copy inserts 0 until 7, with no comment saying why. The expected rows at L1703 stay the same with 0 until 8, since p.b + 1 = 8 is not in bucket4. How about a createPlainTable(bs: Seq[Long] = 0L until 8L) next to createIdentityIdTable? Also, the two new helpers sit between createBucketedIdTable (L740) and its wrapper createBucketedIdTables (L764), which were adjacent before. Could they move after that group?

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

This branch has not been deployed

No deployments
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