Skip to content

Commit d8d72d2

Browse files
committed
[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
### What changes were proposed in this pull request? A storage-partitioned join now pairs a partition transform whose argument is not a bare column, i.e. is an expression or a struct field, only with a transform over the same argument. It also stops dropping such an argument where it used to: - `KeyedShuffleSpec.isExpressionCompatible` pairs a transform with an argument that is not an `Attribute` only with a transform over the same argument shape, i.e. the argument with its one column replaced by a placeholder. So `bucket(4, b + 1)` pairs with `bucket(4, c + 1)` on `b = c`, but not with `bucket(4, x)`. Over the same shape the two compare as two transforms of columns do. Under `allowCompatibleTransforms`, `bucket(8, b + 1)` reduces onto `bucket(4, c + 1)`, since a reduce maps the partition keys and evaluates no argument. When an earlier join reduced the keys, two such transforms pair only when they were reduced through the same pairing. An identity side is never reduced onto such a transform. `TransformExpression` keeps its comparisons. It loses `withReference`, whose last callers now rebuild the transform with `copy`. - `KeyedShuffleSpec.canCreatePartitioning` asks the same question, so such a layout is never the one other children are shuffled onto. `createPartitioning` replaces a transform's argument with the other child's cluster key, which drops whatever surrounds the key. It and the identity arm of `reducersBothWays` now rebuild a transform through one helper that asserts the argument is a column. - `ShuffleSpecCollection.canCreatePartitioning` is now true when any member can, and `EnsureRequirements.pickCoPartitionTarget` offers and ranks only the members that can. So one member that cannot serve no longer rules out a usable sibling. That also covers a bare expression such as `b + 1`, which master already refused. - `GroupPartitionsExec` rebuilds a reduced expression over each member's own argument when it reports the members of a collection. It used to re-target the column only. - `KeyedShuffleSpec.keyPositions` maps an expression with more than one reference, e.g. `bucket(4, b + c)`, to no position instead of failing an assertion (SPARK-59901). - `EnsureRequirements.createKeyedShuffleSpecs` counts only an expression over one column towards `spark.sql.requireAllClusterKeysForCoPartition`. An expression over two columns maps to no position, so a projection would drop it together with its columns. ### Why are the changes needed? With `spark.sql.sources.v2.bucketing.shuffle.enabled` on, this query returns 0 rows instead of 7. `t1(id)` and `t3(x)` are partitioned by `bucket(4, ...)`, `plain(b)` is not partitioned, and each holds the values 0 to 7: ```sql SELECT t1.id, p.b, t3.x FROM t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x ``` 1. The first join shuffles `plain` onto `t1`'s partitioning. `KeyedShuffleSpec.createPartitioning` builds that partitioning from `plain`'s join key, so it is `bucket(4, b + 1)`. This join is correct. 2. The second join clusters `plain` on the bare `b`. `KeyedShuffleSpec.keyPositions` maps a partition expression to a cluster key through its reference, so `bucket(4, b + 1)` counts as a function of `b`, the counterpart of `x`. 3. `isSameFunction` compares only the function name and the bucket count. So `bucket(4, b + 1)` is the same as `t3`'s `bucket(4, x)`, and the join pairs the partitions as they stand. A row with `b = x` sits in bucket `(b + 1) % 4` on one side and in `x % 4` on the other. The comparison ignores which column an argument is, and relies on `keyPositions` to pair the columns up. That is sound only when the argument is the column itself. `bucket(4, b + 1)` is a function of `b`, but not the same function of it as `bucket(4, x)` is of `x`. A scan never reports a transform of an expression or of a struct field, but two planner paths build one. A one-side shuffle does, under `spark.sql.sources.v2.bucketing.shuffle.enabled`. An inner broadcast hash join does too, with every conf at its default. `BroadcastHashJoinExec.expandOutputPartitioning` also reports the streamed side's layout over the build side's join key. The same problem shows up in more shapes, all measured on master: - **A broadcast join.** With every conf at its default, `SELECT /*+ BROADCAST(p) */ ... FROM t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x` returns 0 of 7 rows the same way, through the expanded `bucket(4, b + 1)` layout. - **The reducer path.** Under `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`, a `bucket(8, x)` third table loses the rows the same way. - **A struct field.** A join key `p.s.a` gives `bucket(4, s.a)`, which `keyPositions` pairs through `s`. Joining two such sides on `s` pairs `bucket(4, s.a)` with `bucket(4, s.b)` and returns 0 of 8 rows. Shuffling another side onto such a layout builds `bucket(4, s)`, which fails on executors with a `ClassCastException`. - **Shuffling onto the layout.** In `plain p JOIN t1 ON p.b + 1 = t1.id JOIN xy q ON t1.id = q.x AND p.b = q.y`, the second join shuffles `xy` onto the first join's `bucket(4, b + 1)` layout. That builds `bucket(4, y)` for `xy`, drops the `+ 1`, and returns 0 of 8 rows. - **After a reduce.** When a later join reduces the first join's output from `bucket(4)` onto `bucket(2)`, `GroupPartitionsExec` reports the `bucket(4, b + 1)` member as `bucket(2, b)`. A join on `b` then pairs it with a `bucket(2, x)` scan, or shuffles another side onto it, and returns 0 of 7 rows. The same happens to the bare `b + 1` member of a side shuffled onto an identity-partitioned table. The struct form fails with the `ClassCastException`. - **An identity side reduced onto it.** Under `allowCompatibleTransforms`, a later join can reduce an identity-partitioned side onto `bucket(4, b + 1)`. That evaluates `bucket(4, id + 1)` on the side's partition keys while planning. The keys come out right, but under ANSI a key of `Long.MaxValue` fails the query with `ARITHMETIC_OVERFLOW`, although the query never computes `id + 1`. - **A GROUP BY after a reduce.** With `allowCompatibleTransforms`, `SELECT /*+ BROADCAST(p) */ p.b, max(p.c) FROM ident i JOIN plain2 p ON i.id = p.b + p.c JOIN bucket2 t2 ON i.id = t2.id GROUP BY p.b` returns 4 groups instead of 2. The reduce reports the `p.b + p.c` member as `bucket(2, p.b)`, which serves the `GROUP BY` as it stands. - **A key over two columns.** A join key such as `p.b + p.c` gives `bucket(4, b + c)`, or the bare `b + c` over an identity-partitioned side. Any spec over it fails planning with an `AssertionError` in `keyPositions` (SPARK-59901). AQE validates the plan in `OptimizeSkewedJoin`, so a single join is enough. An `ORDER BY` over such a partition expression had a related problem in another code path. That was SPARK-59905, fixed in #59189. The one-side shuffle of transform expressions came with SPARK-48012, in 4.0.0. So did the broadcast expansion of a keyed layout, since SPARK-49205 made `KeyGroupedPartitioning` an expression. ### Does this PR introduce _any_ user-facing change? Yes, it fixes the wrong results, the crash and the planning failures above. Such a query now shuffles the side instead of pairing it or shuffling onto it. That costs shuffles in three shapes whose plan was already correct: - A widening cast counts as part of the argument's shape. When `t1.id` is `BIGINT` and `p.b` is `INT`, the first join reports `bucket(4, cast(b as bigint))`. A later join `p.b = t3.x` with an `INT` bucketed `t3` then shuffles both sides. Before, a connector whose bucket function has one canonical name for both types paired the two sides as they stood. This happens with every conf at its default, through a broadcast join. Two sides that both carry the cast still pair. - Under `allowCompatibleTransforms`, an identity-partitioned side is no longer reduced onto such a transform, so it is shuffled. - With `spark.sql.sources.v2.bucketing.shuffle.enabled`, a join on two keys whose sides pair only through a transform of an expression gets one more shuffle. Take two broadcast joins, `t1 JOIN /*+ BROADCAST(p) */ p ON t1.id = p.b + 1` and the same over `t2` and `q`, joined on `l.id = r.c AND l.b = r.b`. Each layout covers only one join key, so the storage-partitioned join declines under `requireAllClusterKeysForCoPartition`. Master then keeps both sides, paired on `bucket(4, b + 1)`. This PR cannot shuffle onto that layout, so it shuffles one side onto `bucket(4, id)`. Both plans co-partition on one key only, which `requireAllClusterKeysForCoPartition` is meant to prevent. That is a separate, pre-existing question, filed as SPARK-59971. One plan gets better. A collection with a bare expression member, e.g. the `b + 1` of a side shuffled onto an identity-partitioned table, used to be refused as a layout as a whole. Now its other members can serve, so `ident i JOIN plain p ON i.id = p.b + 1 JOIN xy q ON i.id = q.x AND p.b = q.y` takes 2 shuffles instead of 3. This PR keeps the fix small, since it has to reach the maintenance branches. SPARK-59900 is the follow-up for the rest: shuffling another side onto such a layout, and reducing an identity side onto it. ### How was this patch tested? New tests. All but two of them fail on master: - `ShuffleSpecSuite`, "SPARK-59887: a spec over a transform of an expression pairs only with the same shape". For an expression and a struct field argument, it covers the same shape, another shape and a column, both ways and with itself, the reduce in both directions, the two sides of one reduce, two sides reduced through different pairings, an identity side, `canCreatePartitioning`, and a collection with a usable sibling. - `ShuffleSpecSuite`, "SPARK-59901: an expression over two columns maps to no cluster key". - `KeyGroupedPartitioningSuite`, "SPARK-59887: a side shuffled onto a transform of an expression is not paired as it is". It runs the query above against a `bucket(4)` third table, and with `allowCompatibleTransforms` against a `bucket(8)` one that the join reduces, and asserts two shuffles. - `KeyGroupedPartitioningSuite`, "SPARK-59887: a side is not shuffled onto a transform of an expression". It runs both join orders and asserts that `xy` is shuffled onto `bucket(4, x)` and never onto a bucket of `y`. - `KeyGroupedPartitioningSuite`, "SPARK-59887: a broadcast join's layout over a join key is not paired as it is". It runs with every conf at its default and checks that the first join is a broadcast join. - `KeyGroupedPartitioningSuite`, "SPARK-59887: a layout over a transform of a struct field is not paired nor shuffled onto". Master returns 0 of 8 rows. - `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps what each member's keys are computed from". It covers a bucketed and an identity-partitioned first table, a bucketed and an unpartitioned third table, and the struct form. Master returns 0 of 7 rows. - `KeyGroupedPartitioningSuite`, "SPARK-59887: an identity side is not reduced onto a transform of an expression". It runs with `p.b + 1` over a `Long.MaxValue` key and with `-p.b` over a `Long.MinValue` key. Master fails with the `ARITHMETIC_OVERFLOW` for both. - `KeyGroupedPartitioningSuite`, "SPARK-59887: a reduce keeps a member over two columns". It runs the `GROUP BY` above. Master returns 4 groups instead of 2. - `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns is not paired through one of them". Master fails with the `AssertionError`. - `KeyGroupedPartitioningSuite`, "SPARK-59901: a join key over two columns does not cover its columns". It joins two broadcast joins that report `[b + c, d]` on `(b, c, d)` under `allowKeysSubsetOfPartitionKeys`, and asserts that the join is not paired on `d` alone. Master fails with the `AssertionError`. - `KeyGroupedPartitioningSuite`, "SPARK-59887: two layouts over the same shape of a join key pair as two columns do". It runs the cast case above with every conf at its default. It also runs two one-side shuffles onto `bucket(4, b + 1)` and onto `bucket(4)` or `bucket(8)` of `q.b + 1`, where `allowCompatibleTransforms` reduces the `bucket(8)` one. It asserts that the join adds no shuffle. It passes on master, which pairs and reduces these as well. - `KeyGroupedPartitioningSuite`, "SPARK-59887: the same shape reduced together pairs only through the same pairing". Two legs each shuffle onto `bucket(12, b + 1)`. The first leg reduces with `bucket(8)`, so its keys are buckets of 4. The second leg reduces with `bucket(8)` too, with `bucket(18)`, or not at all, so its keys are buckets of 4, 6 or 12. It asserts 2 shuffles for the first and 4 for the other two. Master returns 0 of 23 rows for the third. - `KeyGroupedPartitioningSuite`, "SPARK-59887: a member that cannot serve does not rank its collection". It pins the ranking: `q` keeps its layout with 3 partitions, as on master, where the collection is refused as a whole. Without the member filter, the collection would rank first with 5 partitions, and `q` would be shuffled onto 2. The SPARK-59905 test from #59189 now also runs its `p.b + p.c` query with AQE on. Master fails that run with the `AssertionError`. The three expression tests each run with the key `p.b + 1` and with `-p.b`. Master returns 0 of 8 rows for `p.b + 1` and 4 of 8 for `-p.b`. The `-p.b` key is the shape that also fails on 4.2, 4.1 and 4.0. On those branches `b + 1` already works, since its literal leaf blocks the pairing. Each change was removed on its own: - without the guard in `isExpressionCompatible`, the one-side shuffle and broadcast pairing tests, the struct field test and the reduce test return wrong rows, and the identity tests fail on the assertion in `rebuiltOver`; - without the shape pairing, the same-shape and same-pairing tests fail on the shuffle count; - with a reduce of such a pair refused, the same-shape test fails on the shuffle count; - with every pair of reduced keys over an expression refused, the same-pairing test fails on the shuffle count; - with the reduced marks ignored, or only checked on both sides, the same-pairing test returns 7 of 23 rows; - without the `canCreatePartitioning` clause, the one-side shuffle pairing, `xy`, struct field, reduce and same-pairing tests fail, all on the assertion in `rebuiltOver`; - with the `forall` back in `ShuffleSpecCollection.canCreatePartitioning`, the `xy` tests fail on the shuffle count; - without the member filter in `pickCoPartitionTarget`, 15 tests fail in `KeyGroupedPartitioningSuite` and `EnsureRequirementsSuite`. The filter now also does what the collection-level check before it did, so 8 of them are older tests. Of this PR's tests, 7 fail, 6 of them on the assertion in `rebuiltOver`. A member over the same shape pairs with itself, so only the filter keeps it from being the layout; - without the `GroupPartitionsExec` change, the reduce test returns 0 of 7 rows, and the two-column `GROUP BY` returns 4 groups; - without the `requireAllClusterKeysForCoPartition` change, the coverage test pairs on `d` alone; - with the `keyPositions` assertion back, both SPARK-59901 tests in `KeyGroupedPartitioningSuite` fail. The first `ShuffleSpecSuite` test also fails for each change to `isExpressionCompatible` above. The `xy` tests pass with the guard removed, since the `canCreatePartitioning` clause and the member filter keep `xy` off that layout. Also ran `TransformExpressionSuite`, `ShuffleSpecSuite`, `DistributionSuite`, the `KeyGroupedPartitioning*` suites, `WriteDistributionAndOrderingSuite`, `PlannerSuite`, `ProjectedOrderingAndPartitioningSuite`, `GroupPartitionsExecSuite`, `EnsureRequirementsSuite`, `BucketedReadWithoutHiveSupportSuite`, the plan stability suites, the `*JoinSuite` suites and `AdaptiveQueryExecSuite`, 1955 tests in all, plus `dev/lint-scala`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5.5) Closes #59165 from peter-toth/SPARK-59887-spj-expression-key-identity. Authored-by: Peter Toth <peter.toth@gmail.com> Signed-off-by: Peter Toth <peter.toth@gmail.com>
1 parent 3a91983 commit d8d72d2

7 files changed

Lines changed: 642 additions & 80 deletions

File tree

‎sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala‎

Lines changed: 0 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -123,16 +123,6 @@ case class TransformExpression(
123123
}
124124
}
125125

126-
/**
127-
* Re-targets this partition transform expression at `attr`. A partition transform expression
128-
* has a single leaf attribute (`KeyedPartitioning.supportsExpressions`), so this replaces that
129-
* attribute and keeps the rest of the expression: any field path above the leaf comes from this
130-
* expression, not from the re-targeted key. Re-targeting is therefore only faithful when the
131-
* source and the target key expressions have the same path shape.
132-
*/
133-
def withReference(attr: Attribute): TransformExpression =
134-
transform { case _: AttributeReference => attr }.asInstanceOf[TransformExpression]
135-
136126
// Return a Reducer for a reducible function on another reducible function
137127
private def reducer(
138128
thisFunction: ReducibleFunction[_, _],

‎sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala‎

Lines changed: 122 additions & 39 deletions
Large diffs are not rendered by default.

‎sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/ShuffleSpecSuite.scala‎

Lines changed: 106 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ package org.apache.spark.sql.catalyst
1919

2020
import org.apache.spark.{SparkFunSuite, SparkUnsupportedOperationException}
2121
import org.apache.spark.sql.catalyst.dsl.expressions._
22-
import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeReference, DirectShufflePartitionID, Expression, TransformExpression}
22+
import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeReference, DirectShufflePartitionID, Expression, GetStructField, TransformExpression}
2323
import org.apache.spark.sql.catalyst.plans.SQLHelper
2424
import org.apache.spark.sql.catalyst.plans.physical._
2525
import org.apache.spark.sql.connector.catalog.functions.{FlipLowBitFunction, Reducer, ReducibleFunction, ScalarFunction}
@@ -761,6 +761,111 @@ class ShuffleSpecSuite extends SparkFunSuite with SQLHelper {
761761
}
762762
}
763763

764+
test("SPARK-59887: a spec over a transform of an expression pairs only with the same shape") {
765+
// A join pairs its two sides up by the column each partition expression references. Over
766+
// `a = b` it would pair `bucket(4, a + 1)` with `bucket(4, b)`, and over `s = t` it would pair
767+
// `bucket(4, s.x)` with `bucket(4, t.y)`. A row with `a = b` sits in a different bucket on each
768+
// side. `bucket(4, a + 1)` and `bucket(4, b + 1)` are the same function of `a` and of `b`, so
769+
// they pair up.
770+
val fn = new FakeBucket
771+
val a = $"a".long
772+
val b = $"b".long
773+
val structType = new StructType().add("x", LongType).add("y", LongType)
774+
val s = $"s".struct(structType)
775+
val t = $"t".struct(structType)
776+
// A spec clustered on the column its one partition expression references.
777+
def keyed(expression: Expression): KeyedShuffleSpec =
778+
KeyedShuffleSpec(
779+
KeyedPartitioning(Seq(expression), Seq(InternalRow(0L), InternalRow(1L))),
780+
ClusteredDistribution(expression.references.toSeq))
781+
def bucket(numBuckets: Int, argument: Expression): TransformExpression =
782+
TransformExpression(fn, Seq(argument), Some(numBuckets))
783+
def spec(numBuckets: Int, argument: Expression): KeyedShuffleSpec =
784+
keyed(bucket(numBuckets, argument))
785+
// Whether the two sides of one reduce, `bucket(8, left)` with `bucket(4, right)`, pair up.
786+
def pairedByReduce(left: Expression, right: Expression): Boolean =
787+
keyed(bucket(8, left).reducedTogetherWith(bucket(4, right)))
788+
.isCompatibleWith(keyed(bucket(4, right).reducedTogetherWith(bucket(8, left))))
789+
val column = spec(4, b)
790+
791+
withSQLConf(
792+
SQLConf.V2_BUCKETING_SHUFFLE_ENABLED.key -> "true",
793+
SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key -> "true",
794+
SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> "false",
795+
SQLConf.V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS.key -> "true") {
796+
assert(spec(4, a).isCompatibleWith(column), "two columns pair up")
797+
assert(spec(4, a).canCreatePartitioning)
798+
assert(spec(8, a).areKeysCompatible(column, allowReduce = true),
799+
"the fixture reduces two columns")
800+
assert(pairedByReduce(a, b), "two columns reduced together share one key space")
801+
assert(keyed(b).areKeysCompatible(spec(4, a), allowReduce = true),
802+
"an identity side reduces onto a column's transform")
803+
804+
Seq(
805+
("an expression", a + 1L, b + 1L, b + 2L),
806+
("a struct field", GetStructField(s, 0), GetStructField(t, 0), GetStructField(t, 1))
807+
).foreach { case (shape, argument, sameShape, otherShape) =>
808+
val over = spec(4, argument)
809+
assert(over.isCompatibleWith(spec(4, sameShape)), s"over $shape, the same shape")
810+
assert(over.isCompatibleWith(over), s"over $shape, with itself")
811+
assert(!over.isCompatibleWith(spec(4, otherShape)), s"over $shape, another shape")
812+
assert(!over.isCompatibleWith(column), s"over $shape, against a column")
813+
assert(!column.isCompatibleWith(over), s"over $shape, the other way round")
814+
assert(!over.canCreatePartitioning, s"over $shape, no layout to shuffle onto")
815+
816+
// Over the same shape, such a pair reduces as two columns do, whichever side reduces. It
817+
// does not reduce onto another shape or a column, and an identity side does not reduce
818+
// onto it.
819+
assert(spec(8, argument).areKeysCompatible(spec(4, sameShape), allowReduce = true),
820+
s"over $shape, reducing the same shape")
821+
assert(spec(4, sameShape).areKeysCompatible(spec(8, argument), allowReduce = true),
822+
s"over $shape, reducing the same shape the other way round")
823+
assert(!spec(8, argument).areKeysCompatible(spec(4, otherShape), allowReduce = true),
824+
s"over $shape, reducing another shape")
825+
assert(!spec(8, argument).areKeysCompatible(column, allowReduce = true), s"over $shape")
826+
assert(!column.areKeysCompatible(spec(8, argument), allowReduce = true), s"over $shape")
827+
assert(!over.areKeysCompatible(spec(8, b), allowReduce = true), s"over $shape")
828+
assert(!keyed(b).areKeysCompatible(over, allowReduce = true),
829+
s"over $shape, reducing an identity side")
830+
// Keys an earlier join reduced are in the space of that reduce. Two transforms reduced
831+
// through the same pairing pair over the same shape, as two columns do. `bucket(12)`
832+
// reduced with `bucket(8)` holds buckets of 4, with `bucket(18)` buckets of 6, and
833+
// unreduced buckets of 12.
834+
assert(pairedByReduce(argument, sameShape), s"over $shape, reduced together")
835+
assert(!pairedByReduce(argument, otherShape),
836+
s"over $shape, another shape reduced together")
837+
val reducedWith8 = keyed(bucket(12, argument).reducedTogetherWith(bucket(8, a)))
838+
assert(!reducedWith8.isCompatibleWith(
839+
keyed(bucket(12, sameShape).reducedTogetherWith(bucket(18, b)))),
840+
s"over $shape, reduced through different pairings")
841+
assert(!reducedWith8.isCompatibleWith(spec(12, sameShape)),
842+
s"over $shape, reduced against unreduced")
843+
assert(!spec(12, sameShape).isCompatibleWith(reducedWith8),
844+
s"over $shape, unreduced against reduced")
845+
846+
// One such member does not keep a sibling from serving as the layout.
847+
assert(ShuffleSpecCollection(Seq(over, spec(4, a))).canCreatePartitioning)
848+
}
849+
}
850+
}
851+
852+
test("SPARK-59901: an expression over two columns maps to no cluster key") {
853+
// A one-side shuffle builds its partition expression over the other side's join key, which can
854+
// reference two columns. No single cluster key stands for it.
855+
val a = $"a".long
856+
val b = $"b".long
857+
Seq(a + b, TransformExpression(new FakeBucket, Seq(a + b), Some(4))).foreach { e =>
858+
val twoColumns = KeyedShuffleSpec(
859+
KeyedPartitioning(Seq(e), Seq(InternalRow(0L), InternalRow(1L))),
860+
ClusteredDistribution(Seq(a, b)))
861+
assert(twoColumns.keyPositions.forall(_.isEmpty), s"$e")
862+
assert(!twoColumns.isCompatibleWith(twoColumns), s"$e")
863+
withSQLConf(SQLConf.V2_BUCKETING_SHUFFLE_ENABLED.key -> "true") {
864+
assert(!twoColumns.canCreatePartitioning, s"$e")
865+
}
866+
}
867+
}
868+
764869
test("createShuffleSpec: a marked narrowing projection yields an unusable spec") {
765870
val a = $"a".int
766871
val b = $"b".int

‎sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala‎

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -111,9 +111,5 @@ class TransformExpressionSuite extends SparkFunSuite {
111111
assert(!left.hasSameReducedKeys(b12), "an unreduced side never shares one")
112112
assert(!b12.hasSameReducedKeys(left))
113113
assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones")
114-
115-
// The marker rides on the expression, so it survives the attribute rewrites a projection and
116-
// `GroupPartitionsExec` apply to a reported partitioning.
117-
assert(left.hasSameReducedKeys(left.withReference(b)))
118114
}
119115
}

‎sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -570,8 +570,8 @@ private[sql] object GroupPartitionsExec {
570570
// can only differ in `expressions`, since they share one `KeyLayout` (enforced by
571571
// `PartitioningCollection`). So the grouping is computed once, and the layout it describes
572572
// is shared by the members below.
573-
// When reducers are applied, the stored reduced expressions are re-targeted at each
574-
// `KeyedPartitioning`'s own key attribute and reported instead of the original ones. Their
573+
// When reducers are applied, the stored reduced expressions are rebuilt over each
574+
// `KeyedPartitioning`'s own argument and reported instead of the original ones. Their
575575
// data types match the reduced partition keys for the identity-vs-transform and
576576
// single-side-transform reducers; for the both-sides-reduce shape no single transform
577577
// describes the keys, so the reduce marks it (see `KeyedShuffleSpec.reducersBothWays`).
@@ -604,8 +604,15 @@ private[sql] object GroupPartitionsExec {
604604
// `reduced` came from the one member `checkKeyGroupCompatible` paired
605605
// this side on, which need not be the member being rewritten. The keys are
606606
// reduced once, from the shared key rows, so `reduced` describes them
607-
// whichever member this is, and only the key attribute is re-targeted.
608-
reduced.withReference(expr.references.head)
607+
// whichever member this is. Only its function says that, so the member's own
608+
// argument takes the place of `reduced`'s. A side shuffled onto this layout
609+
// reports `bucket(8, b + 1)` next to `bucket(8, id)`, or `b + 1` next to an
610+
// identity `id`.
611+
val argument = expr match {
612+
case t: TransformExpression => t.children
613+
case e => Seq(e)
614+
}
615+
reduced.copy(children = argument)
609616
case (expr, None) => expr
610617
}
611618
case None => projectedExpressions

‎sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/EnsureRequirements.scala‎

Lines changed: 19 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -377,25 +377,28 @@ case class EnsureRequirements(
377377
val shouldConsiderMinParallelism = children.zip(specs).forall { case (child, spec) =>
378378
spec.forall(!_.canCreatePartitioning) || child.isInstanceOf[ShuffleExchangeLike]
379379
}
380-
// Choose all the specs that can be used to shuffle other children
381-
val candidateSpecs = children.zip(specs).collect {
382-
case (child, Some(spec)) if spec.canCreatePartitioning &&
383-
(!shouldConsiderMinParallelism ||
384-
child.outputPartitioning.numPartitions >= conf.defaultNumShufflePartitions) =>
385-
child -> spec
380+
// Choose all the children whose spec can be used to shuffle other children, each with the
381+
// members that can serve as the layout. Any other member would build a partitioning its own
382+
// child is not laid out on, and it does not count towards the ranking below either.
383+
val candidateSpecs = children.zip(specs).flatMap {
384+
case (child, Some(spec)) if !shouldConsiderMinParallelism ||
385+
child.outputPartitioning.numPartitions >= conf.defaultNumShufflePartitions =>
386+
val members = spec.flatten.filter(_.canCreatePartitioning)
387+
if (members.nonEmpty) Some(child -> members) else None
388+
case _ => None
386389
}
387390
// Rank on two things at once. A child with no `ShuffleExchangeLike` node comes first, since
388391
// keeping it costs nothing. For instance, if we have:
389392
// A: (No_Exchange, 100) <---> B: (Exchange, 120)
390393
// it's better to pick A and change B to (Exchange, 100) instead of picking B and insert a
391394
// new shuffle for A. Then the best parallelism decides, and for a collection that is the best
392-
// any member offers, since a collection has no count of its own.
395+
// any serving member offers, since a collection has no count of its own.
393396
//
394397
// What the winner contributes is its members, the alternatives the layout is picked from below.
395398
// Empty when no child can serve as the layout.
396-
val bestMembers: Seq[LeafShuffleSpec] = candidateSpecs.maxByOption { case (child, spec) =>
397-
(!child.isInstanceOf[ShuffleExchangeLike], spec.flatten.map(_.numPartitions).max)
398-
}.toSeq.flatMap(_._2.flatten)
399+
val bestMembers: Seq[LeafShuffleSpec] = candidateSpecs.maxByOption { case (child, members) =>
400+
(!child.isInstanceOf[ShuffleExchangeLike], members.map(_.numPartitions).max)
401+
}.toSeq.flatMap(_._2)
399402

400403
// A `ShuffleSpecCollection` answers `isCompatibleWith` if *any* of its members does, so the
401404
// winner alone does not say which member the sides agreed on. The projection pushed into a
@@ -408,7 +411,7 @@ case class EnsureRequirements(
408411
val pairings: Seq[(LeafShuffleSpec, Seq[Option[LeafShuffleSpec]])] = bestMembers.map { member =>
409412
member -> childMembers.map(_.find(member.isCompatibleWith))
410413
}
411-
// Which children the winner reaches at all, over all of its members.
414+
// Which children the winner reaches at all, over all of its serving members.
412415
val reached: Seq[Boolean] = pairings.map(_._2).transpose.map(_.exists(_.isDefined))
413416
// The member picked has to pair with every child the winner reaches. A child it does not
414417
// reach, and a child that is not in the decision, are shuffled whichever member wins, so
@@ -1133,9 +1136,11 @@ case class EnsureRequirements(
11331136
// the skew of joining on keys that are coarser than the join keys. Key order and duplicated
11341137
// cluster keys don't matter.
11351138
def allClusterKeysCovered: Boolean =
1136-
// The single-column invariant in KeyedPartitioning.supportsExpressions guarantees one
1137-
// attribute per partition expression.
1138-
distribution.allClusterKeysAmong(partitioning.expressions.flatMap(_.references))
1139+
// Only an expression over a single column covers that column. One over several, e.g.
1140+
// `b + c`, maps to no position (`KeyedShuffleSpec.keyPositions`). The spec turns it away
1141+
// or projects it away, so its columns are not covered.
1142+
distribution.allClusterKeysAmong(
1143+
partitioning.expressions.filter(_.references.size == 1).flatMap(_.references))
11391144

11401145
// The coverage requirement is a comparison of expressions, while `keysMaySatisfy` can end in
11411146
// a projection of the partition keys, so the cheap question is asked first. The requirement

0 commit comments

Comments
 (0)