From 791e616ea12ef0f4705fbec952b38c1ec7d0c88c Mon Sep 17 00:00:00 2001 From: Peter Toth Date: Tue, 6 Oct 2026 09:11:08 +0200 Subject: [PATCH 1/4] [SPARK-59995][SQL] Let `TransformExpression` generate the code of its bound function `TransformExpression.doGenCode` now generates the code of the bound function call that `eval` runs. It still throws when there is no such call. The k-way merge of `GroupPartitionsExec` compares rows with generated code. An ordering with a partition transform, reported by the source or derived from the partition keys, made the task fail with `Cannot generate code for expression`. --- .../expressions/TransformExpression.scala | 9 ++- .../TransformExpressionSuite.scala | 59 ++++++++++++++- .../KeyGroupedPartitioningSuite.scala | 72 +++++++++++++++++++ 3 files changed, 137 insertions(+), 3 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala index bfa94914f4ae9..3118e573c3ff0 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala @@ -186,8 +186,15 @@ case class TransformExpression( case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this) } + // `resolvedFunction` is not a child, so a tree walk does not see it. For a function with only + // `produceResult` it falls back to `eval` on the input row, which whole-stage codegen does not + // provide (see `CodegenFallback.generate`). No whole-stage operator generates code for a + // transform today. override protected def doGenCode(ctx: CodegenContext, ev: ExprCode): ExprCode = - throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this) + resolvedFunction match { + case Some(fn) => fn.genCode(ctx) + case None => throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this) + } } object TransformExpression { diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala index d333b41441618..d553e4a26983c 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala @@ -17,11 +17,14 @@ package org.apache.spark.sql.catalyst.expressions -import org.apache.spark.SparkFunSuite +import org.apache.spark.{SparkException, SparkFunSuite} +import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.sql.catalyst.expressions.codegen.CodegenContext +import org.apache.spark.sql.catalyst.expressions.objects.{Invoke, StaticInvoke} import org.apache.spark.sql.connector.catalog.functions.{BoundFunction, ScalarFunction} import org.apache.spark.sql.types.{DataType, IntegerType} -class TransformExpressionSuite extends SparkFunSuite { +class TransformExpressionSuite extends SparkFunSuite with ExpressionEvalHelper { /** * A bound function with a stable canonical name and no `equals` of its own. A plain class, NOT a @@ -112,4 +115,56 @@ class TransformExpressionSuite extends SparkFunSuite { assert(!b12.hasSameReducedKeys(left)) assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones") } + + test("SPARK-59995: the generated code computes what eval does") { + // `checkEvaluation` runs both the interpreted and the generated path. Each function resolves to + // one of the calls `V2ExpressionUtils.resolveScalarFunction` builds. + val input = BoundReference(0, IntegerType, nullable = true) + Seq( + InvokeBucketFunction -> classOf[Invoke], + new StaticInvokeBucketFunction -> classOf[StaticInvoke], + ProduceResultBucketFunction -> classOf[ApplyFunctionExpression] + ).foreach { case (fn, call) => + val transform = bucket(fn, input) + assert(call.isInstance(transform.resolvedFunction.get), transform.resolvedFunction) + checkEvaluation(transform, 3, create_row(7)) + checkEvaluation(transform, 2, create_row(6)) + } + } + + test("SPARK-59995: a transform whose keys a join reduced generates no code") { + // The call no longer computes such keys, so it must not run in the generated path either. + val input = BoundReference(0, IntegerType, nullable = true) + val reduced = bucket(InvokeBucketFunction, input, 12) + .reducedTogetherWith(bucket(InvokeBucketFunction, input, 8)) + checkError( + exception = intercept[SparkException](reduced.genCode(new CodegenContext)), + condition = "INTERNAL_ERROR", + parameters = Map("message" -> s"Cannot generate code for expression: $reduced")) + } +} + +/** `bucket` over integers. The functions below differ only in how Spark calls them. */ +private trait IntBucketFunction extends ScalarFunction[Int] { + override def inputTypes(): Array[DataType] = Array(IntegerType, IntegerType) + override def resultType(): DataType = IntegerType + override def name(): String = "bucket" +} + +/** Called through the magic `invoke` method, which has generated code. */ +private object InvokeBucketFunction extends IntBucketFunction { + def invoke(numBuckets: Int, value: Int): Int = Math.floorMod(value, numBuckets) +} + +/** Called through a static `invoke` method, as a Java function would be. */ +private class StaticInvokeBucketFunction extends IntBucketFunction + +private object StaticInvokeBucketFunction { + def invoke(numBuckets: Int, value: Int): Int = Math.floorMod(value, numBuckets) +} + +/** Called through `produceResult`, which falls back to `eval`. */ +private object ProduceResultBucketFunction extends IntBucketFunction { + override def produceResult(input: InternalRow): Int = + Math.floorMod(input.getInt(1), input.getInt(0)) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala index 0b037ed6829b4..66a8babf775f8 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala @@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite } } + /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces them. */ + private def insertItemsInTwoYears(): Unit = sql(s"INSERT INTO testcat.ns.$items VALUES " + + "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " + + "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " + + "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))") + + /** + * Joins `items` with `purchases` on `joinCondition` and checks that the plan k-way merges over an + * ordering with a partition transform. + */ + private def checkKWayMergeOverTransform(joinCondition: String): Unit = { + val df = sql( + s""" + |${selectWithMergeJoinHint("i", "p")} + |i.id, i.name, i.arrive_time + |FROM testcat.ns.$items i JOIN testcat.ns.$purchases p ON $joinCondition + |""".stripMargin) + checkAnswer(df, Seq( + Row(1, "aa", Timestamp.valueOf("2021-01-01 00:00:00")), + Row(1, "ab", Timestamp.valueOf("2022-01-01 00:00:00")), + Row(2, "bb", Timestamp.valueOf("2021-01-01 00:00:00")))) + val merging = collectAllGroupPartitions(df.queryExecution.executedPlan) + .filter(_.enableSortedMerge) + assert(merging.length == 1, "expected one k-way merge") + assert(merging.head.child.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]), + "expected a partition transform in the merge's ordering") + assert(merging.head.execute().isInstanceOf[SortedMergeCoalescedRDD[_]]) + } + + test("SPARK-59995: k-way merge over a reported ordering with a partition transform") { + // The join on (id, name) needs the merge, since the key ordering on id alone is not enough. + // The merge's ordering is [id, name, years(arrive_time)], so generating its comparator + // generates code for the transform. + val itemOrdering = Array( + sort(FieldReference("id"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST), + sort(FieldReference("name"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST), + sort(years("arrive_time"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST)) + createTable(items, itemsColumns, Array(identity("id")), itemOrdering) + insertItemsInTwoYears() + val namedPurchasesColumns = Array( + Column.create("item_id", LongType), + Column.create("name", StringType)) + createTable(purchases, namedPurchasesColumns, Array(identity("item_id"))) + sql(s"INSERT INTO testcat.ns.$purchases VALUES (1, 'aa'), (1, 'ab'), (2, 'bb')") + + withSQLConf( + SQLConf.REQUIRE_ALL_CLUSTER_KEYS_FOR_CO_PARTITION.key -> "false", + SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true") { + checkKWayMergeOverTransform("p.item_id = i.id AND p.name = i.name") + } + } + + test("SPARK-59995: k-way merge over an ordering derived from a partition transform key") { + // The scan reports no ordering, so it derives [id, years(arrive_time)] from its keys. The join + // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced + // partitions keep no ordering on id unless they are merged. + createTable(items, itemsColumns, Array(identity("id"), years("arrive_time"))) + insertItemsInTwoYears() + createTable(purchases, purchasesColumns, Array(identity("item_id"))) + sql(s"INSERT INTO testcat.ns.$purchases VALUES " + + "(1, 10.0, cast('2021-01-01' as timestamp)), " + + "(2, 20.0, cast('2021-01-01' as timestamp))") + + withSQLConf( + SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true", + SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true", + SQLConf.V2_BUCKETING_PRESERVE_KEY_ORDERING_ON_COALESCE_ENABLED.key -> "false") { + checkKWayMergeOverTransform("p.item_id = i.id") + } + } + test("SPARK-56549: k-way merge enabled only when parent requires ordering") { // Dynamic gate: with the config enabled, k-way merge must be activated only when the parent // actually requires ordering (SMJ), and must stay off when the parent does not (hash join). From 913f317a3456886235f816ec714a0d13c4d993b9 Mon Sep 17 00:00:00 2001 From: Peter Toth Date: Wed, 7 Oct 2026 10:40:01 +0200 Subject: [PATCH 2/4] Drop sort orders that hold a partition transform from a V2 scan's output ordering Reverts the doGenCode delegation. No operator requires an ordering over a transform, so the k-way merge no longer compares rows by one. The k-way merge comparator, now LazyRowOrdering, builds its comparator with RowOrdering.create, so it falls back to interpreted evaluation when code generation fails. --- .../expressions/TransformExpression.scala | 9 +-- .../TransformExpressionSuite.scala | 59 +---------------- .../v2/DataSourceV2ScanExecBase.scala | 18 +++-- .../datasources/v2/GroupPartitionsExec.scala | 23 ++++--- .../KeyGroupedPartitioningSuite.scala | 66 +++++++++++++------ .../WriteDistributionAndOrderingSuite.scala | 6 +- .../v2/GroupPartitionsExecSuite.scala | 28 ++++++-- 7 files changed, 100 insertions(+), 109 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala index 3118e573c3ff0..bfa94914f4ae9 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala @@ -186,15 +186,8 @@ case class TransformExpression( case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this) } - // `resolvedFunction` is not a child, so a tree walk does not see it. For a function with only - // `produceResult` it falls back to `eval` on the input row, which whole-stage codegen does not - // provide (see `CodegenFallback.generate`). No whole-stage operator generates code for a - // transform today. override protected def doGenCode(ctx: CodegenContext, ev: ExprCode): ExprCode = - resolvedFunction match { - case Some(fn) => fn.genCode(ctx) - case None => throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this) - } + throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this) } object TransformExpression { diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala index d553e4a26983c..d333b41441618 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala @@ -17,14 +17,11 @@ package org.apache.spark.sql.catalyst.expressions -import org.apache.spark.{SparkException, SparkFunSuite} -import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.expressions.codegen.CodegenContext -import org.apache.spark.sql.catalyst.expressions.objects.{Invoke, StaticInvoke} +import org.apache.spark.SparkFunSuite import org.apache.spark.sql.connector.catalog.functions.{BoundFunction, ScalarFunction} import org.apache.spark.sql.types.{DataType, IntegerType} -class TransformExpressionSuite extends SparkFunSuite with ExpressionEvalHelper { +class TransformExpressionSuite extends SparkFunSuite { /** * A bound function with a stable canonical name and no `equals` of its own. A plain class, NOT a @@ -115,56 +112,4 @@ class TransformExpressionSuite extends SparkFunSuite with ExpressionEvalHelper { assert(!b12.hasSameReducedKeys(left)) assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones") } - - test("SPARK-59995: the generated code computes what eval does") { - // `checkEvaluation` runs both the interpreted and the generated path. Each function resolves to - // one of the calls `V2ExpressionUtils.resolveScalarFunction` builds. - val input = BoundReference(0, IntegerType, nullable = true) - Seq( - InvokeBucketFunction -> classOf[Invoke], - new StaticInvokeBucketFunction -> classOf[StaticInvoke], - ProduceResultBucketFunction -> classOf[ApplyFunctionExpression] - ).foreach { case (fn, call) => - val transform = bucket(fn, input) - assert(call.isInstance(transform.resolvedFunction.get), transform.resolvedFunction) - checkEvaluation(transform, 3, create_row(7)) - checkEvaluation(transform, 2, create_row(6)) - } - } - - test("SPARK-59995: a transform whose keys a join reduced generates no code") { - // The call no longer computes such keys, so it must not run in the generated path either. - val input = BoundReference(0, IntegerType, nullable = true) - val reduced = bucket(InvokeBucketFunction, input, 12) - .reducedTogetherWith(bucket(InvokeBucketFunction, input, 8)) - checkError( - exception = intercept[SparkException](reduced.genCode(new CodegenContext)), - condition = "INTERNAL_ERROR", - parameters = Map("message" -> s"Cannot generate code for expression: $reduced")) - } -} - -/** `bucket` over integers. The functions below differ only in how Spark calls them. */ -private trait IntBucketFunction extends ScalarFunction[Int] { - override def inputTypes(): Array[DataType] = Array(IntegerType, IntegerType) - override def resultType(): DataType = IntegerType - override def name(): String = "bucket" -} - -/** Called through the magic `invoke` method, which has generated code. */ -private object InvokeBucketFunction extends IntBucketFunction { - def invoke(numBuckets: Int, value: Int): Int = Math.floorMod(value, numBuckets) -} - -/** Called through a static `invoke` method, as a Java function would be. */ -private class StaticInvokeBucketFunction extends IntBucketFunction - -private object StaticInvokeBucketFunction { - def invoke(numBuckets: Int, value: Int): Int = Math.floorMod(value, numBuckets) -} - -/** Called through `produceResult`, which falls back to `eval`. */ -private object ProduceResultBucketFunction extends IntBucketFunction { - override def produceResult(input: InternalRow): Int = - Math.floorMod(input.getInt(1), input.getInt(0)) } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala index 49bee1e7961c7..4ac78db59bf4c 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala @@ -19,7 +19,7 @@ package org.apache.spark.sql.execution.datasources.v2 import org.apache.spark.rdd.RDD import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, SortOrder} +import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, SortOrder, TransformExpression} import org.apache.spark.sql.catalyst.plans.physical import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning import org.apache.spark.sql.catalyst.util.truncatedString @@ -144,19 +144,29 @@ trait DataSourceV2ScanExecBase * is a `KeyedPartitioning` and `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` * is on, each partition contains rows where the key expressions evaluate to a single constant * value, so the data is trivially sorted by those expressions within the partition. + * + * Either way, a sort order that holds a partition transform is dropped, even on a partition key. + * In a reported ordering it also ends the leading run, like a sort order over a pruned column. + * This loses nothing today, since no operator requires an ordering over a transform. The write + * path sorts by the transform's function call instead. Dropping it is also a safe way to handle + * a transform Spark cannot evaluate. Keeping it would only add comparisons nobody uses. A sort + * order that Spark derives from a transform key is even constant within each partition. Revisit + * this if an ordering over a transform becomes a real requirement. */ override def outputOrdering: Seq[SortOrder] = { + def holdsTransform(e: Expression): Boolean = e.exists(_.isInstanceOf[TransformExpression]) (ordering, outputPartitioning) match { case (Some(o), p) => - val (prefix, rest) = o.span(_.references.subsetOf(outputSet)) + val (prefix, rest) = + o.span(order => order.references.subsetOf(outputSet) && !holdsTransform(order.child)) p match { case k: KeyedPartitioning if rest.nonEmpty => - val keyExprs = ExpressionSet(k.expressions) + val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform)) prefix ++ rest.filter(order => keyExprs.contains(order.child)) case _ => prefix } case (_, k: KeyedPartitioning) if conf.v2BucketingPartitionKeyOrderingEnabled => - k.expressions.map(SortOrder(_, Ascending)) + k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending)) case _ => Seq.empty } } diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala index 2cdc6b9d410b1..1c1afc8dd051e 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala @@ -24,7 +24,6 @@ import org.apache.spark.{Partition, SparkException} import org.apache.spark.rdd.{CoalescedRDD, PartitionCoalescer, PartitionGroup, RDD, SortedMergeCoalescedRDD} import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.catalyst.expressions._ -import org.apache.spark.sql.catalyst.expressions.codegen.GenerateOrdering import org.apache.spark.sql.catalyst.plans.QueryPlan import org.apache.spark.sql.catalyst.plans.physical.{IdentityReducer, KeyedPartitioning, KeyLayout, KeyReducer, Partitioning, PartitioningCollection, UnknownPartitioning} import org.apache.spark.sql.catalyst.util.{truncatedString, InternalRowComparableWrapper} @@ -267,8 +266,8 @@ case class GroupPartitionsExec( } /** - * The ordering used by the k-way merge in [[SortedMergeCoalescedRDD]]. The generated comparator - * ([[GenerateOrdering]]) only needs each [[SortOrder]]'s sort key (child, direction, null + * The ordering used by the k-way merge in [[SortedMergeCoalescedRDD]]. The comparator that + * [[RowOrdering]] builds only needs each [[SortOrder]]'s sort key (child, direction, null * ordering), so `sameOrderExpressions` -- planner-only metadata that would otherwise be * serialized with the RDD in every task -- is dropped. */ @@ -335,7 +334,7 @@ case class GroupPartitionsExec( sparkContext.emptyRDD } else if (usesSortedMerge) { val partitionCoalescer = new GroupedPartitionCoalescer(groupedPartitions.map(_._2)) - val rowOrdering = new LazyCodeGenOrdering(kWayMergeOrdering, child.output) + val rowOrdering = new LazyRowOrdering(kWayMergeOrdering, child.output) new SortedMergeCoalescedRDD[InternalRow]( child.execute(), groupedPartitions.size, @@ -798,15 +797,15 @@ class GroupedPartitionCoalescer( } /** - * A serializable [[Ordering]] for [[InternalRow]] that generates code-compiled comparison logic - * lazily on first use. The [[SortOrder]] expressions and output schema are serialized with the - * RDD; the generated comparator is rebuilt on the executor on first comparison via - * [[GenerateOrdering]]. + * A serializable [[Ordering]] for [[InternalRow]]. The [[SortOrder]] expressions and output schema + * are serialized with the RDD. The comparator is built on the executor, at the first comparison. + * [[RowOrdering]] builds it, as for `SortExec`. It tries generated code first and falls back to + * interpreted evaluation when that fails. */ -private class LazyCodeGenOrdering( +private class LazyRowOrdering( sortOrders: Seq[SortOrder], schema: Seq[Attribute]) extends Ordering[InternalRow] with Serializable { - @transient private lazy val generated: Ordering[InternalRow] = - GenerateOrdering.generate(sortOrders, schema) - override def compare(x: InternalRow, y: InternalRow): Int = generated.compare(x, y) + @transient private lazy val ordering: Ordering[InternalRow] = + RowOrdering.create(sortOrders, schema) + override def compare(x: InternalRow, y: InternalRow): Int = ordering.compare(x, y) } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala index 66a8babf775f8..5de3c38878917 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala @@ -6571,17 +6571,44 @@ class KeyGroupedPartitioningSuite } } - /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces them. */ - private def insertItemsInTwoYears(): Unit = sql(s"INSERT INTO testcat.ns.$items VALUES " + - "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " + - "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " + - "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))") + test("SPARK-59995: a scan's output ordering drops sort orders that hold a partition transform") { + // `t1` reports a transform in the middle, so the leading run of its ordering stops there. + // `t2` reports a transform key first. Each split holds a single key, so the sort order on `id` + // after it still holds. `t3` reports no ordering, so the scan derives one from its keys. In + // each case a sort on `id` right above the scan needs no `SortExec`. + val table1 = "transform_order_t1" + val table2 = "transform_order_t2" + val table3 = "transform_order_t3" + def asc(expr: Expression): SortOrder = + sort(expr, SortDirection.ASCENDING, NullOrdering.NULLS_FIRST) + createTable(table1, columns, Array(identity("id")), + Array(asc(FieldReference("id")), asc(years("ts")), asc(FieldReference("data")))) + createTable(table2, columns, Array(years("ts"), identity("id")), + Array(asc(years("ts")), asc(FieldReference("id")))) + createTable(table3, columns, Array(days("ts"), identity("id"))) + Seq(table1, table2, table3).foreach { table => + sql(s"INSERT INTO testcat.ns.$table VALUES (1, 'aa', cast('2020-01-01' as timestamp))") + val df = sql(s"SELECT id, data, ts FROM testcat.ns.$table") + val scan = collectScans(df.queryExecution.executedPlan).head + assert(scan.outputOrdering.map(_.child.sql) === Seq("id"), table) + val sorted = df.sortWithinPartitions("id").queryExecution.executedPlan + assert(collect(sorted) { case s: SortExec => s }.isEmpty, table) + } + } /** - * Joins `items` with `purchases` on `joinCondition` and checks that the plan k-way merges over an - * ordering with a partition transform. + * Inserts three items, two of them with `id` 1 on their own splits, so grouping the splits by + * `id` coalesces them. Then joins `items` with `purchases` on `joinCondition` and checks that + * the plan k-way merges over `expectedOrdering`, the columns the scan's ordering keeps before + * its transform. */ - private def checkKWayMergeOverTransform(joinCondition: String): Unit = { + private def checkKWayMergeBeforeTransform( + joinCondition: String, + expectedOrdering: Seq[String]): Unit = { + sql(s"INSERT INTO testcat.ns.$items VALUES " + + "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " + + "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " + + "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))") val df = sql( s""" |${selectWithMergeJoinHint("i", "p")} @@ -6595,21 +6622,18 @@ class KeyGroupedPartitioningSuite val merging = collectAllGroupPartitions(df.queryExecution.executedPlan) .filter(_.enableSortedMerge) assert(merging.length == 1, "expected one k-way merge") - assert(merging.head.child.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]), - "expected a partition transform in the merge's ordering") + assert(merging.head.child.outputOrdering.map(_.child.sql) == expectedOrdering) assert(merging.head.execute().isInstanceOf[SortedMergeCoalescedRDD[_]]) } test("SPARK-59995: k-way merge over a reported ordering with a partition transform") { // The join on (id, name) needs the merge, since the key ordering on id alone is not enough. - // The merge's ordering is [id, name, years(arrive_time)], so generating its comparator - // generates code for the transform. + // The scan reports [id, name, years(arrive_time)] and keeps [id, name]. val itemOrdering = Array( sort(FieldReference("id"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST), sort(FieldReference("name"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST), sort(years("arrive_time"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST)) createTable(items, itemsColumns, Array(identity("id")), itemOrdering) - insertItemsInTwoYears() val namedPurchasesColumns = Array( Column.create("item_id", LongType), Column.create("name", StringType)) @@ -6619,16 +6643,17 @@ class KeyGroupedPartitioningSuite withSQLConf( SQLConf.REQUIRE_ALL_CLUSTER_KEYS_FOR_CO_PARTITION.key -> "false", SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true") { - checkKWayMergeOverTransform("p.item_id = i.id AND p.name = i.name") + checkKWayMergeBeforeTransform("p.item_id = i.id AND p.name = i.name", Seq("id", "name")) } } test("SPARK-59995: k-way merge over an ordering derived from a partition transform key") { - // The scan reports no ordering, so it derives [id, years(arrive_time)] from its keys. The join - // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced - // partitions keep no ordering on id unless they are merged. - createTable(items, itemsColumns, Array(identity("id"), years("arrive_time"))) - insertItemsInTwoYears() + // The scan reports no ordering, so it derives [id, days(arrive_time)] from its keys and keeps + // [id]. The merge must not call the transform's function. Spark cannot call this `days` at + // all, since `DaysFunction` implements neither `invoke` nor `produceResult`. The join on id + // projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced partitions + // keep no ordering on id unless they are merged. + createTable(items, itemsColumns, Array(identity("id"), days("arrive_time"))) createTable(purchases, purchasesColumns, Array(identity("item_id"))) sql(s"INSERT INTO testcat.ns.$purchases VALUES " + "(1, 10.0, cast('2021-01-01' as timestamp)), " + @@ -6636,10 +6661,9 @@ class KeyGroupedPartitioningSuite withSQLConf( SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true", - SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true", SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true", SQLConf.V2_BUCKETING_PRESERVE_KEY_ORDERING_ON_COALESCE_ENABLED.key -> "false") { - checkKWayMergeOverTransform("p.item_id = i.id") + checkKWayMergeBeforeTransform("p.item_id = i.id", Seq("id")) } } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala index ec1b34b8c2101..2a2c0a7bb5564 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/WriteDistributionAndOrderingSuite.scala @@ -1545,10 +1545,12 @@ class WriteDistributionAndOrderingSuite extends DistributionAndOrderingSuiteBase val df = sql("SELECT id, data FROM testcat.ns1.test_table") val scans = collect(df.queryExecution.executedPlan) { case s: BatchScanExec => s } assert(scans.size === 1) - val ordering = scans.head.outputOrdering + val ordering = scans.head.ordering.getOrElse(Seq.empty) assert(ordering.nonEmpty, - "scan should report non-empty outputOrdering via SupportsReportOrdering") + "scan should keep the ordering reported via SupportsReportOrdering") assert(ordering.head.child.isInstanceOf[TransformExpression], "bucket-based sort order should resolve to a TransformExpression") + // SPARK-59995: the output ordering drops the sort order over the transform. + assert(scans.head.outputOrdering.isEmpty) } } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala index fb9afa9cf3f3f..f305da1dd864c 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala @@ -17,17 +17,19 @@ package org.apache.spark.sql.execution.datasources.v2 -import org.apache.spark.{SparkContext, SparkException} +import org.apache.spark.{SparkConf, SparkContext, SparkException} import org.apache.spark.rdd.RDD +import org.apache.spark.serializer.JavaSerializer import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, AttributeReference, SortOrder, TransformExpression} +import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, AttributeReference, CodegenObjectFactoryMode, SortOrder, TransformExpression} import org.apache.spark.sql.catalyst.plans.physical.{ClusteredDistribution, KeyedPartitioning, KeyReducer, Partitioning, PartitioningCollection, UnknownPartitioning} +import org.apache.spark.sql.catalyst.util.DateTimeTestUtils.date import org.apache.spark.sql.catalyst.util.InternalRowComparableWrapper -import org.apache.spark.sql.connector.catalog.functions.{BucketFunction, BucketReducer, DaysFunctionWithToYearsReducerWithLongResult, DaysToYearsReducerWithLongResult, Reducer, YearsFunctionWithToYearsReducerWithLongResult} +import org.apache.spark.sql.connector.catalog.functions.{BucketFunction, BucketReducer, DaysFunctionWithToYearsReducerWithLongResult, DaysToYearsReducerWithLongResult, Reducer, YearsFunction, YearsFunctionWithToYearsReducerWithLongResult} import org.apache.spark.sql.execution.{DummySparkPlan, LeafExecNode, SafeForKWayMerge} import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.test.SharedSparkSession -import org.apache.spark.sql.types.{DataType, IntegerType, LongType} +import org.apache.spark.sql.types.{DataType, IntegerType, LongType, TimestampType} import org.apache.spark.sql.vectorized.ColumnarBatch class GroupPartitionsExecSuite extends SharedSparkSession { @@ -144,7 +146,7 @@ class GroupPartitionsExecSuite extends SharedSparkSession { test("SPARK-58324: k-way merge ordering drops sameOrderExpressions") { // The child ordering carries sameOrderExpressions (planner metadata). The k-way merge // comparator only needs the sort key, so kWayMergeOrdering keeps child/direction/nullOrdering - // but drops sameOrderExpressions, so LazyCodeGenOrdering does not serialize them with the RDD. + // but drops sameOrderExpressions, so LazyRowOrdering does not serialize them with the RDD. val childOrdering = Seq(SortOrder(exprA, Ascending, Seq(exprB, exprC))) val child = DummySparkPlan( outputPartitioning = KeyedPartitioning(Seq(exprA), Seq(row(1), row(2), row(1))), @@ -157,6 +159,22 @@ class GroupPartitionsExecSuite extends SharedSparkSession { assert(merged.forall(_.sameOrderExpressions.isEmpty)) } + test("SPARK-59995: the k-way merge ordering falls back to interpreted evaluation") { + // A `TransformExpression` generates no code, while `eval` calls its function. The ordering is + // serialized before any comparison, as the RDD ships it. Test sessions run `CODEGEN_ONLY`, so + // this sets the production default, `FALLBACK`. + val ts = AttributeReference("ts", TimestampType)() + val ordering = new LazyRowOrdering( + Seq(SortOrder(TransformExpression(YearsFunction, Seq(ts)), Ascending)), Seq(ts)) + val serializer = new JavaSerializer(new SparkConf()).newInstance() + val shipped = serializer.deserialize[LazyRowOrdering](serializer.serialize(ordering)) + withSQLConf(SQLConf.CODEGEN_FACTORY_MODE.key -> CodegenObjectFactoryMode.FALLBACK.toString) { + assert(shipped.compare(InternalRow(date(2022)), InternalRow(date(2021, 6))) > 0) + // Two timestamps in the same year compare equal, so the year is what is compared. + assert(shipped.compare(InternalRow(date(2021)), InternalRow(date(2021, 6))) === 0) + } + } + test("SPARK-56241: coalescing without reducers keeps key-expression orders from child") { // Key 1 appears on partitions 0 and 2, causing coalescing. val partitionKeys = Seq(row(1), row(2), row(1)) From 6eca5ef5e41493d5fe497342e651b51ac952ac08 Mon Sep 17 00:00:00 2001 From: Peter Toth Date: Wed, 7 Oct 2026 19:37:18 +0200 Subject: [PATCH 3/4] Address review comments Pin partitionKeyOrdering in the tests for the maintenance branches, keep the transform's column in the scan output on purpose, check the codegen premise of the fallback test, and say in the docs that the derived ordering leaves partition transforms out. --- docs/sql-migration-guide.md | 2 +- docs/sql-performance-tuning.md | 2 +- .../catalog/functions/BoundFunction.java | 2 +- .../read/SupportsReportOrdering.java | 4 ++ .../apache/spark/sql/internal/SQLConf.scala | 3 +- .../v2/DataSourceV2ScanExecBase.scala | 10 ++--- .../datasources/v2/GroupPartitionsExec.scala | 8 ++-- .../KeyGroupedPartitioningSuite.scala | 37 ++++++++++++------- .../v2/GroupPartitionsExecSuite.scala | 20 ++++++++-- 9 files changed, 56 insertions(+), 32 deletions(-) diff --git a/docs/sql-migration-guide.md b/docs/sql-migration-guide.md index bebd6ab0e5562..eb12b8d847d42 100644 --- a/docs/sql-migration-guide.md +++ b/docs/sql-migration-guide.md @@ -42,7 +42,7 @@ license: | - Since Spark 4.4, when `array_repeat` or `array_insert` is asked to build an array larger than the maximum supported array length, generated code raises the same error as interpreted evaluation. `array_repeat` now fails with `COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER` instead of the internal error `_LEGACY_ERROR_TEMP_2176`, and `array_insert` fails with `COLLECTION_SIZE_LIMIT_EXCEEDED.FUNCTION` instead of `COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER`, which named a `count` parameter that `array_insert` does not have. Both functions raise an error under exactly the same conditions as before; only the reported error condition changes. - Since Spark 4.4, `/*+ ... */` inside a bracketed comment is parsed as a nested comment, so its closing `*/` no longer closes the outer comment. For example, `/* note /*+ x */ SELECT 1` previously returned `1`, but now raises `UNCLOSED_BRACKETED_COMMENT`. Close the outer comment explicitly, for example `/* note /*+ x */ */ SELECT 1`. - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partition.filter.enabled` defaults to `true`. In a storage-partitioned join, the partition key groups pushed down to both sides may now be narrowed to those that can produce output for the join type, instead of always taking the union of the two sides' groups: an inner or semi join may keep only the groups present on both sides, a join that keeps or tests every left row (left outer, left anti, left single, existence) keeps the left side's groups, a right outer join keeps the right side's, and a full outer join is unaffected; a CROSS join carrying an equality condition is treated like an inner join. The narrowing is skipped when either side's partitioning may contain unknown partition keys, as after a shuffle on one side, so such a join keeps the full union. A group that is dropped is not scanned at all, so a query may read fewer files. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partition.filter.enabled` to `false`. -- Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate whose grouping is exactly the partition key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. +- Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate that groups by those key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. - Since Spark 4.4, `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` defaults to `true`. When `GroupPartitionsExec` merges several input partitions that share one partition key value into a single output partition, it now keeps sort orders over the partition key expressions in the ordering it reports, since those expressions are constant within the merged partition. Orders over other columns are still dropped, as the concatenation invalidates them, and a join that reduced the partition keys onto a common key space (see `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`) reports no order, since the merged partitions then share only the reduced key. This can remove a `Sort` downstream of the merge. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` to `false`. ## Upgrading from Spark SQL 4.2 to 4.3 diff --git a/docs/sql-performance-tuning.md b/docs/sql-performance-tuning.md index 232c47d16b13a..762b9063d8148 100644 --- a/docs/sql-performance-tuning.md +++ b/docs/sql-performance-tuning.md @@ -708,7 +708,7 @@ The following SQL properties enable Storage Partition Join in different join que spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled true - When enabled, Spark derives the output ordering of a V2 scan from its partition key expressions, if the source reports a keyed partitioning but no explicit ordering, or an ordering that Spark ignores because it references a column that cannot be resolved. All rows of such a partition share one key value, so the partition is trivially sorted by those expressions, and a sort Spark would otherwise add becomes unnecessary. This config requires spark.sql.sources.v2.bucketing.enabled to be true. + When enabled, Spark derives the output ordering of a V2 scan from its partition key expressions, if the source reports a keyed partitioning but no explicit ordering, or an ordering that Spark ignores because it references a column that cannot be resolved. All rows of such a partition share one key value, so the partition is trivially sorted by those expressions, and a sort Spark would otherwise add becomes unnecessary. Partition transforms such as days(ts) or bucket(8, id) are left out of the ordering. This config requires spark.sql.sources.v2.bucketing.enabled to be true. 4.2.0 diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java index a5b729a6fa33b..7a55cd3a4fc70 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java @@ -118,7 +118,7 @@ default String canonicalName() { *
    *
  • accepting a valid query that selects and groups by the same scalar function call
  • *
  • keeping a union's keyed partitioning
  • - *
  • retaining a reported ordering that matches the partitioning
  • + *
  • merging scans that report an ordering over the same transform
  • *
  • reusing identical bucketed scans
  • *
  • recognizing two identical subplans
  • *
diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java index cf1a0532d3c07..eebcb1f2e7599 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java @@ -39,6 +39,10 @@ public interface SupportsReportOrdering extends Scan { * Spark resolves the column references in these sort orders against the table columns, * including columns pruned from the scan output. If any of them cannot be resolved, Spark * ignores the whole reported ordering and logs a warning. + *

+ * Spark does not use a sort order over a transform such as {@code days(ts)}, whether or not it is + * a partition key. It drops that sort order and every one after it, except the sort orders on + * partition keys that are not transforms. */ SortOrder[] outputOrdering(); } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index 5801d8e399c45..b40a55a61ddc7 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -2592,7 +2592,8 @@ object SQLConf { "a V2 data source that reports a KeyedPartitioning but does not report explicit ordering " + "via SupportsReportOrdering, or reports one that Spark ignores because it references a " + "column that cannot be resolved. Within a single partition all rows share the same key " + - s"value, so the data is trivially sorted by those expressions. Requires " + + "value, so the data is trivially sorted by those expressions. Partition transforms such " + + "as `days(ts)` or `bucket(8, id)` are left out of the ordering. Requires " + s"${V2_BUCKETING_ENABLED.key} to be enabled.") .version("4.2.0") .withBindingPolicy(ConfigBindingPolicy.SESSION) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala index 4ac78db59bf4c..62ce0156c3f2f 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala @@ -147,11 +147,11 @@ trait DataSourceV2ScanExecBase * * Either way, a sort order that holds a partition transform is dropped, even on a partition key. * In a reported ordering it also ends the leading run, like a sort order over a pruned column. - * This loses nothing today, since no operator requires an ordering over a transform. The write - * path sorts by the transform's function call instead. Dropping it is also a safe way to handle - * a transform Spark cannot evaluate. Keeping it would only add comparisons nobody uses. A sort - * order that Spark derives from a transform key is even constant within each partition. Revisit - * this if an ordering over a transform becomes a real requirement. + * Dropping it loses nothing today when Spark can call the transform's function, since no + * operator then requires an ordering over the transform. The write path sorts by the function + * call instead. Keeping such a sort order would add comparisons nobody uses. Dropping it also + * handles a transform Spark cannot evaluate. Revisit this if an ordering over a transform + * becomes a real requirement. */ override def outputOrdering: Seq[SortOrder] = { def holdsTransform(e: Expression): Boolean = e.exists(_.isInstanceOf[TransformExpression]) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala index 1c1afc8dd051e..c9ee6780a3281 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala @@ -388,11 +388,9 @@ case class GroupPartitionsExec( outputPartitioning match { case p: Partitioning with Expression if reducers.isEmpty && conf.v2BucketingPreserveKeyOrderingOnCoalesceEnabled => - // Without reducers all merged partitions share the same original key value, so the key - // expressions remain constant within the output partition. The child's outputOrdering - // should already be in sync with the partitioning (either reported by the source or - // derived from it in DataSourceV2ScanExecBase), so we only need to keep the sort orders - // whose expression is a partition key expression -- all others are lost by concatenation. + // Without reducers all merged partitions share the same original key value, so the sort + // orders on key expressions still hold. The transform keys match nothing here, since + // `DataSourceV2ScanExecBase.outputOrdering` drops every sort order over a transform. val keyedPartitionings = p.collect { case k: KeyedPartitioning => k } val keyExprs = ExpressionSet(keyedPartitionings.flatMap(_.expressions)) child.outputOrdering.filter(order => keyExprs.contains(order.child)) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala index 5de3c38878917..cd53ec3d9e1b6 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala @@ -6575,24 +6575,28 @@ class KeyGroupedPartitioningSuite // `t1` reports a transform in the middle, so the leading run of its ordering stops there. // `t2` reports a transform key first. Each split holds a single key, so the sort order on `id` // after it still holds. `t3` reports no ordering, so the scan derives one from its keys. In - // each case a sort on `id` right above the scan needs no `SortExec`. + // each case a sort on `id` right above the scan needs no `SortExec`. The query selects `ts`, so + // the transform's column stays in the scan output. Otherwise SPARK-59899's pruned-column + // handling would give the same orderings even without this fix. val table1 = "transform_order_t1" val table2 = "transform_order_t2" val table3 = "transform_order_t3" - def asc(expr: Expression): SortOrder = - sort(expr, SortDirection.ASCENDING, NullOrdering.NULLS_FIRST) + def asc(expr: Expression): SortOrder = sort(expr, SortDirection.ASCENDING) createTable(table1, columns, Array(identity("id")), Array(asc(FieldReference("id")), asc(years("ts")), asc(FieldReference("data")))) createTable(table2, columns, Array(years("ts"), identity("id")), Array(asc(years("ts")), asc(FieldReference("id")))) createTable(table3, columns, Array(days("ts"), identity("id"))) - Seq(table1, table2, table3).foreach { table => - sql(s"INSERT INTO testcat.ns.$table VALUES (1, 'aa', cast('2020-01-01' as timestamp))") - val df = sql(s"SELECT id, data, ts FROM testcat.ns.$table") - val scan = collectScans(df.queryExecution.executedPlan).head - assert(scan.outputOrdering.map(_.child.sql) === Seq("id"), table) - val sorted = df.sortWithinPartitions("id").queryExecution.executedPlan - assert(collect(sorted) { case s: SortExec => s }.isEmpty, table) + withSQLConf(SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true") { + Seq(table1, table2, table3).foreach { table => + sql(s"INSERT INTO testcat.ns.$table VALUES (1, 'aa', cast('2020-01-01' as timestamp))") + val df = sql(s"SELECT id, data, ts FROM testcat.ns.$table") + val scan = collectScans(df.queryExecution.executedPlan).head + assert(scan.output.exists(_.name == "ts"), s"test setup: ts stays in the output of $table") + assert(scan.outputOrdering.map(_.child.sql) === Seq("id"), table) + val sorted = df.sortWithinPartitions("id").queryExecution.executedPlan + assert(collect(sorted) { case s: SortExec => s }.isEmpty, table) + } } } @@ -6600,7 +6604,9 @@ class KeyGroupedPartitioningSuite * Inserts three items, two of them with `id` 1 on their own splits, so grouping the splits by * `id` coalesces them. Then joins `items` with `purchases` on `joinCondition` and checks that * the plan k-way merges over `expectedOrdering`, the columns the scan's ordering keeps before - * its transform. + * its transform. The query selects `arrive_time`, so the transform's column stays in the scan + * output. Otherwise SPARK-59899's pruned-column handling would drop the transform even without + * this fix. */ private def checkKWayMergeBeforeTransform( joinCondition: String, @@ -6622,6 +6628,8 @@ class KeyGroupedPartitioningSuite val merging = collectAllGroupPartitions(df.queryExecution.executedPlan) .filter(_.enableSortedMerge) assert(merging.length == 1, "expected one k-way merge") + assert(merging.head.child.output.exists(_.name == "arrive_time"), + "test setup: arrive_time stays in the scan output") assert(merging.head.child.outputOrdering.map(_.child.sql) == expectedOrdering) assert(merging.head.execute().isInstanceOf[SortedMergeCoalescedRDD[_]]) } @@ -6630,9 +6638,9 @@ class KeyGroupedPartitioningSuite // The join on (id, name) needs the merge, since the key ordering on id alone is not enough. // The scan reports [id, name, years(arrive_time)] and keeps [id, name]. val itemOrdering = Array( - sort(FieldReference("id"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST), - sort(FieldReference("name"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST), - sort(years("arrive_time"), SortDirection.ASCENDING, NullOrdering.NULLS_FIRST)) + sort(FieldReference("id"), SortDirection.ASCENDING), + sort(FieldReference("name"), SortDirection.ASCENDING), + sort(years("arrive_time"), SortDirection.ASCENDING)) createTable(items, itemsColumns, Array(identity("id")), itemOrdering) val namedPurchasesColumns = Array( Column.create("item_id", LongType), @@ -6661,6 +6669,7 @@ class KeyGroupedPartitioningSuite withSQLConf( SQLConf.V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.key -> "true", + SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true", SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> "true", SQLConf.V2_BUCKETING_PRESERVE_KEY_ORDERING_ON_COALESCE_ENABLED.key -> "false") { checkKWayMergeBeforeTransform("p.item_id = i.id", Seq("id")) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala index f305da1dd864c..83eaccbb1f4d3 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala @@ -160,14 +160,26 @@ class GroupPartitionsExecSuite extends SharedSparkSession { } test("SPARK-59995: the k-way merge ordering falls back to interpreted evaluation") { - // A `TransformExpression` generates no code, while `eval` calls its function. The ordering is - // serialized before any comparison, as the RDD ships it. Test sessions run `CODEGEN_ONLY`, so - // this sets the production default, `FALLBACK`. + // The scan drops every sort order over a transform, so none reaches the merge. Here a transform + // only stands in for a sort key without generated code. `TransformExpression` generates none, + // while `eval` calls its function. The ordering is serialized before any comparison, as the + // RDD ships it. val ts = AttributeReference("ts", TimestampType)() val ordering = new LazyRowOrdering( Seq(SortOrder(TransformExpression(YearsFunction, Seq(ts)), Ascending)), Seq(ts)) val serializer = new JavaSerializer(new SparkConf()).newInstance() - val shipped = serializer.deserialize[LazyRowOrdering](serializer.serialize(ordering)) + def ship(): LazyRowOrdering = + serializer.deserialize[LazyRowOrdering](serializer.serialize(ordering)) + // Under `CODEGEN_ONLY` the comparator cannot be built, which is the premise of this test. + val error = withSQLConf( + SQLConf.CODEGEN_FACTORY_MODE.key -> CodegenObjectFactoryMode.CODEGEN_ONLY.toString) { + intercept[SparkException] { + ship().compare(InternalRow(date(2022)), InternalRow(date(2021))) + } + } + assert(error.getMessage.contains("Cannot generate code for expression")) + // `FALLBACK` is the production default. + val shipped = ship() withSQLConf(SQLConf.CODEGEN_FACTORY_MODE.key -> CodegenObjectFactoryMode.FALLBACK.toString) { assert(shipped.compare(InternalRow(date(2022)), InternalRow(date(2021, 6))) > 0) // Two timestamps in the same year compare equal, so the year is what is compared. From f77f37e79e690794530cbe5c076d2931361652fd Mon Sep 17 00:00:00 2001 From: Peter Toth Date: Thu, 8 Oct 2026 10:02:56 +0200 Subject: [PATCH 4/4] Address more review comments Narrow the SupportsReportOrdering and BoundFunction Javadocs, say in the preserveKeyOrderingOnCoalesce docs that transform keys are left out, and check the sort aggregate in the scan ordering test. --- docs/sql-migration-guide.md | 4 ++-- docs/sql-performance-tuning.md | 2 +- .../catalog/functions/BoundFunction.java | 2 +- .../read/SupportsReportOrdering.java | 7 +++--- .../apache/spark/sql/internal/SQLConf.scala | 7 +++--- .../KeyGroupedPartitioningSuite.scala | 24 ++++++++++++------- 6 files changed, 27 insertions(+), 19 deletions(-) diff --git a/docs/sql-migration-guide.md b/docs/sql-migration-guide.md index eb12b8d847d42..f9c7bf4771227 100644 --- a/docs/sql-migration-guide.md +++ b/docs/sql-migration-guide.md @@ -42,8 +42,8 @@ license: | - Since Spark 4.4, when `array_repeat` or `array_insert` is asked to build an array larger than the maximum supported array length, generated code raises the same error as interpreted evaluation. `array_repeat` now fails with `COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER` instead of the internal error `_LEGACY_ERROR_TEMP_2176`, and `array_insert` fails with `COLLECTION_SIZE_LIMIT_EXCEEDED.FUNCTION` instead of `COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER`, which named a `count` parameter that `array_insert` does not have. Both functions raise an error under exactly the same conditions as before; only the reported error condition changes. - Since Spark 4.4, `/*+ ... */` inside a bracketed comment is parsed as a nested comment, so its closing `*/` no longer closes the outer comment. For example, `/* note /*+ x */ SELECT 1` previously returned `1`, but now raises `UNCLOSED_BRACKETED_COMMENT`. Close the outer comment explicitly, for example `/* note /*+ x */ */ SELECT 1`. - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partition.filter.enabled` defaults to `true`. In a storage-partitioned join, the partition key groups pushed down to both sides may now be narrowed to those that can produce output for the join type, instead of always taking the union of the two sides' groups: an inner or semi join may keep only the groups present on both sides, a join that keeps or tests every left row (left outer, left anti, left single, existence) keeps the left side's groups, a right outer join keeps the right side's, and a full outer join is unaffected; a CROSS join carrying an equality condition is treated like an inner join. The narrowing is skipped when either side's partitioning may contain unknown partition keys, as after a shuffle on one side, so such a join keeps the full union. A group that is dropped is not scanned at all, so a query may read fewer files. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partition.filter.enabled` to `false`. -- Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate that groups by those key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. -- Since Spark 4.4, `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` defaults to `true`. When `GroupPartitionsExec` merges several input partitions that share one partition key value into a single output partition, it now keeps sort orders over the partition key expressions in the ordering it reports, since those expressions are constant within the merged partition. Orders over other columns are still dropped, as the concatenation invalidates them, and a join that reduced the partition keys onto a common key space (see `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`) reports no order, since the merged partitions then share only the reduced key. This can remove a `Sort` downstream of the merge. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` to `false`. +- Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. This lets Spark drop a `Sort` it would otherwise place above the scan. An aggregate that groups by a leading run of the other key expressions, in key order, may also be planned as a sort aggregate rather than a hash aggregate. This is because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. +- Since Spark 4.4, `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` defaults to `true`. When `GroupPartitionsExec` merges several input partitions that share one partition key value into a single output partition, it now keeps sort orders over the partition key expressions in the ordering it reports, since those expressions are constant within the merged partition. Sort orders over partition transforms such as `days(ts)` or `bucket(8, id)` are left out. Orders over other columns are still dropped, as the concatenation invalidates them, and a join that reduced the partition keys onto a common key space (see `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`) reports no order, since the merged partitions then share only the reduced key. This can remove a `Sort` downstream of the merge. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` to `false`. ## Upgrading from Spark SQL 4.2 to 4.3 diff --git a/docs/sql-performance-tuning.md b/docs/sql-performance-tuning.md index 762b9063d8148..fa0c2e05663ed 100644 --- a/docs/sql-performance-tuning.md +++ b/docs/sql-performance-tuning.md @@ -716,7 +716,7 @@ The following SQL properties enable Storage Partition Join in different join que spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled true - When enabled, GroupPartitionsExec reports sort orders over partition key expressions after coalescing several input partitions into one. The merged partitions share the same partition key value, so these orders still hold, while orders over other columns are lost by the concatenation. No order is reported when the join reduced the partition keys onto a common key space (see spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled), because the merged partitions then share only the reduced key. This config requires spark.sql.sources.v2.bucketing.enabled to be true. + When enabled, GroupPartitionsExec reports sort orders over partition key expressions after coalescing several input partitions into one. The merged partitions share the same partition key value, so these orders still hold, while orders over other columns are lost by the concatenation. Sort orders over partition transforms such as days(ts) or bucket(8, id) are left out. No order is reported when the join reduced the partition keys onto a common key space (see spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled), because the merged partitions then share only the reduced key. This config requires spark.sql.sources.v2.bucketing.enabled to be true. 4.2.0 diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java index 7a55cd3a4fc70..d227acaa73b6c 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java @@ -118,7 +118,7 @@ default String canonicalName() { *

    *
  • accepting a valid query that selects and groups by the same scalar function call
  • *
  • keeping a union's keyed partitioning
  • - *
  • merging scans that report an ordering over the same transform
  • + *
  • merging scans that report a partitioning or an ordering over the same transform
  • *
  • reusing identical bucketed scans
  • *
  • recognizing two identical subplans
  • *
diff --git a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java index eebcb1f2e7599..95cf76e7f66f3 100644 --- a/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java +++ b/sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java @@ -40,9 +40,10 @@ public interface SupportsReportOrdering extends Scan { * including columns pruned from the scan output. If any of them cannot be resolved, Spark * ignores the whole reported ordering and logs a warning. *

- * Spark does not use a sort order over a transform such as {@code days(ts)}, whether or not it is - * a partition key. It drops that sort order and every one after it, except the sort orders on - * partition keys that are not transforms. + * Spark currently does not rely on a sort order over a transform such as {@code days(ts)} to + * avoid a sort, whether or not it is a partition key. Nor does it rely on the sort orders after + * it. The exception is a sort order on a partition key that is not a transform, when Spark uses + * the reported {@code KeyGroupedPartitioning}. */ SortOrder[] outputOrdering(); } diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala index b40a55a61ddc7..e8f1dc40e69c2 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala @@ -2605,9 +2605,10 @@ object SQLConf { .doc("When enabled, Spark preserves sort orders over partition key expressions when " + "GroupPartitionsExec coalesces multiple input partitions into one output partition. " + "Because all merged partitions share the same partition key value, sort orders over " + - "those key expressions remain valid after the merge. This applies to both key-derived " + - "ordering (from SupportsReportOrdering) and ordering derived from " + - s"${V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key}. Requires " + + "those key expressions remain valid after the merge. This applies to both the ordering " + + "reported via SupportsReportOrdering and the ordering derived from " + + s"${V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key}. Sort orders over partition " + + "transforms such as `days(ts)` or `bucket(8, id)` are left out. Requires " + s"${V2_BUCKETING_ENABLED.key} to be enabled.") .version("4.2.0") .withBindingPolicy(ConfigBindingPolicy.SESSION) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala index cd53ec3d9e1b6..914ac90fd5500 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala @@ -6575,9 +6575,10 @@ class KeyGroupedPartitioningSuite // `t1` reports a transform in the middle, so the leading run of its ordering stops there. // `t2` reports a transform key first. Each split holds a single key, so the sort order on `id` // after it still holds. `t3` reports no ordering, so the scan derives one from its keys. In - // each case a sort on `id` right above the scan needs no `SortExec`. The query selects `ts`, so - // the transform's column stays in the scan output. Otherwise SPARK-59899's pruned-column - // handling would give the same orderings even without this fix. + // each case a sort on `id` right above the scan needs no `SortExec`. A `GROUP BY id` plans a + // `SortAggregateExec`. The query selects `ts`, so the transform's column stays in the scan + // output. Otherwise SPARK-59899's pruned-column handling would give the same orderings even + // without this fix. val table1 = "transform_order_t1" val table2 = "transform_order_t2" val table3 = "transform_order_t3" @@ -6587,7 +6588,9 @@ class KeyGroupedPartitioningSuite createTable(table2, columns, Array(years("ts"), identity("id")), Array(asc(years("ts")), asc(FieldReference("id")))) createTable(table3, columns, Array(days("ts"), identity("id"))) - withSQLConf(SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true") { + withSQLConf( + SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true", + SQLConf.REPLACE_HASH_WITH_SORT_AGG_ENABLED.key -> "true") { Seq(table1, table2, table3).foreach { table => sql(s"INSERT INTO testcat.ns.$table VALUES (1, 'aa', cast('2020-01-01' as timestamp))") val df = sql(s"SELECT id, data, ts FROM testcat.ns.$table") @@ -6596,6 +6599,9 @@ class KeyGroupedPartitioningSuite assert(scan.outputOrdering.map(_.child.sql) === Seq("id"), table) val sorted = df.sortWithinPartitions("id").queryExecution.executedPlan assert(collect(sorted) { case s: SortExec => s }.isEmpty, table) + val aggregated = sql(s"SELECT id, max(ts) FROM testcat.ns.$table GROUP BY id") + .queryExecution.executedPlan + assert(collect(aggregated) { case a: SortAggregateExec => a }.nonEmpty, table) } } } @@ -6656,11 +6662,11 @@ class KeyGroupedPartitioningSuite } test("SPARK-59995: k-way merge over an ordering derived from a partition transform key") { - // The scan reports no ordering, so it derives [id, days(arrive_time)] from its keys and keeps - // [id]. The merge must not call the transform's function. Spark cannot call this `days` at - // all, since `DaysFunction` implements neither `invoke` nor `produceResult`. The join on id - // projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced partitions - // keep no ordering on id unless they are merged. + // The scan reports no ordering, so it derives [id] from its keys, leaving out + // days(arrive_time). The merge must not call the transform's function. Spark cannot call this + // `days` at all, since `DaysFunction` implements neither `invoke` nor `produceResult`. The join + // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced + // partitions keep no ordering on id unless they are merged. createTable(items, itemsColumns, Array(identity("id"), days("arrive_time"))) createTable(purchases, purchasesColumns, Array(identity("item_id"))) sql(s"INSERT INTO testcat.ns.$purchases VALUES " +