Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,11 @@ public interface SupportsReportOrdering extends Scan {

/**
* Returns the order in each partition of this data source scan.
* <p>
* 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();
}
Original file line number Diff line number Diff line change
Expand Up @@ -2241,7 +2241,8 @@ object SQLConf {
.doc("When enabled, Spark derives output ordering from the partition key expressions of " +
"a V2 data source that reports a KeyedPartitioning but does not report explicit ordering " +
"via SupportsReportOrdering. 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)
Expand All @@ -2253,9 +2254,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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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, RowOrdering, SortOrder}
import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, RowOrdering, 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
Expand Down Expand Up @@ -116,19 +116,29 @@ trait DataSourceV2ScanExecBase
* `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
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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, KeyReducer, Partitioning, UnknownPartitioning}
import org.apache.spark.sql.catalyst.util.{truncatedString, InternalRowComparableWrapper}
Expand Down Expand Up @@ -347,8 +346,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.
*/
Expand All @@ -360,7 +359,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,
Expand Down Expand Up @@ -412,11 +411,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))
Expand Down Expand Up @@ -489,15 +486,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)
}
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ import org.apache.spark.sql.execution.{
SortExec,
SparkPlan,
UnionExec}
import org.apache.spark.sql.execution.aggregate.SortAggregateExec
import org.apache.spark.sql.execution.datasources.v2.{BatchScanExec, DataSourceV2ScanRelation, GroupPartitionsExec}
import org.apache.spark.sql.execution.exchange.{ShuffleExchangeExec, ShuffleExchangeLike, ValidateRequirements}
import org.apache.spark.sql.execution.joins.{BroadcastHashJoinExec, ShuffledHashJoinExec, SortMergeJoinExec}
Expand Down Expand Up @@ -5511,6 +5512,119 @@ class KeyGroupedPartitioningSuite extends DistributionAndOrderingSuiteBase with
}
}

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. Before Spark 4.4, a join on a
// subset of the partition keys also needs requireAllClusterKeysForCoPartition off.
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.REQUIRE_ALL_CLUSTER_KEYS_FOR_CO_PARTITION.key -> "false",
SQLConf.V2_BUCKETING_ALLOW_JOIN_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).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
Loading