peter-toth commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4210536223
##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,102 @@ 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`.
+ 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)
+ }
+ }
+
+ /**
+ * 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 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.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,
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)
+ 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, 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)), " +
+ "(2, 20.0, cast('2021-01-01' as timestamp))")
+
+ withSQLConf(
Review Comment:
Done in 6eca5ef5e41. The shape test and the derived merge test now set
`partitionKeyOrdering` explicitly. The 4.2 backport will use the conf's 4.2
name for the subset-keys setting.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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
Review Comment:
Thanks, reworded in 6eca5ef5e41: "The write path sorts by the transform's
function call instead, when Spark can call the function."
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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))
Review Comment:
Agreed. I filed SPARK-60043 for it.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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))
Review Comment:
Done in 6eca5ef5e41, in the conf doc, the tuning guide, the 4.4 migration
note and the `SupportsReportOrdering` Javadoc.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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 =>
Review Comment:
Added both to the user-facing section. I measured the sort aggregate:
`SELECT id, max(ts) FROM t GROUP BY id` over the keys `[years(ts), id]` plans
its partial aggregate as a `SortAggregateExec` with this PR, and as a
`HashAggregateExec` without it. One correction: `preserveOrderingOnCoalesce` is
off by default on every branch. So the description says that the failure needs
it, and that the derived case also needs `partitionKeyOrdering`, which defaults
to true only on master and branch-4.x.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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))
Review Comment:
Dropped in 6eca5ef5e41, but I put a different bullet in its place.
`PlanMerger` compares the logical `ordering`, which keeps the transform, so
merging two scans that report an ordering over the same transform still depends
on a semantic `equals`.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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.
Review Comment:
Reworded in 6eca5ef5e41.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -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
Review Comment:
Shortened in 6eca5ef5e41. "only" and "even constant" are gone, and the
write-path sentence now has the qualifier from item 14.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]