diff --git a/docs/sql-migration-guide.md b/docs/sql-migration-guide.md index bebd6ab0e5562..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 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.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 232c47d16b13a..fa0c2e05663ed 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.enabledspark.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.
spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabledGroupPartitionsExec 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.
+ * 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 c93968e16a7f6..d3211aba4986f 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 @@ -2655,7 +2655,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) @@ -2667,9 +2668,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/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..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 @@ -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. + * 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]) (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 9f4ceb93dab5d..7147721e0197b 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, REPLICATED_FOR_JOIN, UngroupingOrigin, UnknownPartitioning} import org.apache.spark.sql.catalyst.util.{truncatedString, InternalRowComparableWrapper} @@ -277,8 +276,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. */ @@ -345,7 +344,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, @@ -399,11 +398,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)) @@ -820,15 +817,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 2b767c452939b..5b333989c1365 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 @@ -6573,6 +6573,117 @@ class KeyGroupedPartitioningSuite } } + 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`. 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" + 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"))) + 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") + 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) + 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) + } + } + } + + /** + * 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. 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, + 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")} + |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.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[_]]) + } + + 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 scan reports [id, name, years(arrive_time)] and keeps [id, name]. + val itemOrdering = Array( + 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), + 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") { + 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] 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 " + + "(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") { + checkKWayMergeBeforeTransform("p.item_id = i.id", Seq("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). 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 7bab293edf364..c122a7f6c795b 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, OrderedDistribution, Partitioning, PartitioningCollection, REPLICATED_FOR_JOIN, SPLIT_FOR_JOIN, 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,34 @@ class GroupPartitionsExecSuite extends SharedSparkSession { assert(merged.forall(_.sameOrderExpressions.isEmpty)) } + test("SPARK-59995: the k-way merge ordering falls back to interpreted evaluation") { + // 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() + 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. + 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))