peter-toth commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4217072858
##########
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.
+ * <p>
+ * Spark does not use a sort order over a transform such as {@code
days(ts)}, whether or not it is
Review Comment:
Thanks, reworded in f77f37e79e6: "Spark currently does not rely on a sort
order over a transform such as `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 `KeyGroupedPartitioning`." I used "rely on ... to avoid a sort" rather
than "when planning", since `PlanMerger` still compares these sort orders
(SPARK-60083).
##########
sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/BoundFunction.java:
##########
@@ -118,7 +118,7 @@ default String canonicalName() {
* <ul>
* <li>accepting a valid query that selects and groups by the same scalar
function call</li>
* <li>keeping a union's keyed partitioning</li>
- * <li>retaining a reported ordering that matches the partitioning</li>
+ * <li>merging scans that report an ordering over the same transform</li>
Review Comment:
Done in f77f37e79e6.
##########
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.
+ * 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])
Review Comment:
Agreed. I filed SPARK-60083 for it.
##########
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`.
Review Comment:
Done in f77f37e79e6, in all three places. The conf doc also called the
reported ordering "key-derived", which I fixed while there.
##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,111 @@ class KeyGroupedPartitioningSuite
}
}
+ test("SPARK-59995: a scan's output ordering drops sort orders that hold a
partition transform") {
Review Comment:
Added in f77f37e79e6. The shape test now checks that `GROUP BY id` plans a
`SortAggregateExec` for each table, with `replaceHashWithSortAgg` set
explicitly.
##########
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`.
Review Comment:
I'd leave it to the PR description. The change only removes work, and this
fix goes back to branch-4.3 and branch-4.2, where a 4.4 migration note would
not apply.
##########
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`.
Review Comment:
Done in f77f37e79e6. The transform sentence comes first, and the aggregate
one now says "groups by a leading run of the other key expressions, in key
order".
##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,111 @@ 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`. 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") {
+ 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)
+ }
+ }
+ }
+
+ /**
+ * 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, days(arrive_time)]
from its keys and keeps
Review Comment:
Fixed in f77f37e79e6.
--
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]