dongjoon-hyun commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4210790392
##########
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:
Minor: this new public Javadoc says more than holds.
- The sort orders on partition keys after a transform are kept only when
Spark keeps the reported `KeyGroupedPartitioning`, i.e. with
`spark.sql.sources.v2.bucketing.enabled` on and every split implementing
`HasPartitionKey`. Otherwise they are dropped too
(DataSourceV2ScanExecBase.scala:166).
- "Spark does not use a sort order over a transform" holds for physical
planning only. `PlanMerger.combineRequiredOrdering` and
`mergeDegradesReporting` (PlanMerger.scala:967-1001) still compare the reported
ordering with the transform in it, which is what the new `BoundFunction` bullet
says.
- SPARK-60043 plans to change "every one after it", so this public interface
doc would have to change again.
Could this be less committal, e.g. "Spark currently does not exploit a sort
order over a transform such as {@code days(ts)} when planning, nor the sort
orders after it, except those on partition keys that are not transforms when
Spark uses the reported {@code KeyGroupedPartitioning}"?
##########
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:
Minor: the scan merge depends on `equals` for a reported partitioning over a
transform too, not only for an ordering.
`PlanMerger.combineRequiredKeyGroupedPartitioning` (PlanMerger.scala:955-961)
compares the two inputs' reported keys by canonical equality, and the
SPARK-58769 test in MergeSubplansSuite ("decline the merge when the reported
transform is not comparable") pins exactly that for `bucket(4, a)`. That is the
more common case, and after this PR it is the one that still matters for the
physical plan.
Could this bullet say "merging scans that report a partitioning or an
ordering over the same transform"?
##########
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:
Nit: the sort-aggregate change that the PR description lists as user-facing
has no test: `SELECT id, max(ts) FROM t GROUP BY id` over the keys `[years(ts),
id]` now plans a partial `SortAggregateExec`. A later change to
`ReplaceHashWithSortAgg` or to this ordering could flip it silently.
Could this test also check, for `t2`, that `GROUP BY id` plans a
`SortAggregateExec`, with `spark.sql.execution.replaceHashWithSortAgg` set
explicitly for the backports?
##########
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:
Minor: the transform sort orders are dropped only here, in the physical
`outputOrdering`. The logical `DataSourceV2ScanRelation.ordering` keeps them,
so `PlanMerger` (`combineRequiredOrdering` and `mergeDegradesReporting`,
PlanMerger.scala:967-1001) and `DataSourceV2ScanRelation.doCanonicalize`
(DataSourceV2Relation.scala:414-418) still gate scan merging and subplan
deduplication on sort orders that no plan uses any more. For example, take a
`SCAN_MERGING` table that reports `[years(ts)]` and whose `years` function does
not override `equals`. Two scalar subqueries over it bind `years` separately,
so neither ordering satisfies the other and the merge is declined (unless
`spark.sql.optimizer.mergeSubplans.dsv2ScanMerge.orderingDegradation.enabled`
is on), although both physical scans now report `[]`. The new `BoundFunction`
bullet documents this cost rather than removing it.
Could we file a follow-up JIRA to leave the transform sort orders out of
those logical comparisons as well?
##########
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:
Nit: no migration note covers the reported-ordering half of this change,
which applies under any conf. A source that reports `[years(ts), id]` over the
keys `[years(ts), id]` now gets `[id]`, which can drop a `Sort` and turn a hash
aggregate into a sort aggregate, as the PR description says. This note covers
only the derived ordering, and its "To restore the previous behavior" sentence
does not apply to the reported one. It only removes sorts, so this is low
priority.
Could we add a sentence for it, or would you prefer to leave it to the PR
description?
##########
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:
Minor: like the `partitionKeyOrdering` docs in item 16, the
`preserveKeyOrderingOnCoalesce` docs now overstate what is kept: this note
("keeps sort orders over the partition key expressions"), the conf doc
(SQLConf.scala:2605-2610) and docs/sql-performance-tuning.md:719. Sort orders
over transform keys no longer reach `GroupPartitionsExec`. With keys
`[days(ts)]` and a reported `[days(ts)]`, a coalescing `GroupPartitionsExec`
reported `[days(ts)]` before this PR and reports nothing now.
Could we add the same "partition transforms such as `days(ts)` or `bucket(8,
id)` are left out" to those three places?
##########
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:
Nit: with this PR the scan never derives `[id, days(arrive_time)]`.
`k.expressions.filterNot(holdsTransform)` (DataSourceV2ScanExecBase.scala:169)
leaves the transform key out of the derivation itself, so nothing is cut down
to `[id]`.
Could this say "so it derives [id] from its keys, leaving out
days(arrive_time)"?
##########
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:
Nit: "An aggregate that groups by those key expressions" now comes before
the sentence that leaves the transforms out, so it reads as if grouping by a
transform key could switch. Also, "exactly" is gone, although
`ReplaceHashWithSortAgg` matches the ordering by prefix. With keys `[a, id]`,
`GROUP BY id` or `GROUP BY id, a` still plans a hash aggregate, and with keys
`[days(ts)]` no aggregate can switch.
Could the transform sentence come first, and the aggregate one say "groups
by a leading run of those key columns, in key order"?
--
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]