dongjoon-hyun commented on code in PR #59257:
URL: https://github.com/apache/spark/pull/59257#discussion_r4235361065


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -155,19 +161,21 @@ trait DataSourceV2ScanExecBase
    */
   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(order => order.references.subsetOf(outputSet) && 
!holdsTransform(order.child))
-        p match {
-          case k: KeyedPartitioning if rest.nonEmpty =>
-            val keyExprs = 
ExpressionSet(k.expressions.filterNot(holdsTransform))
-            prefix ++ rest.filter(order => keyExprs.contains(order.child))
-          case _ => prefix
-        }
-      case (_, k: KeyedPartitioning) if 
conf.v2BucketingPartitionKeyOrderingEnabled =>
+    val reported = ordering.map { o =>
+      val (prefix, rest) =
+        o.span(order => order.references.subsetOf(outputSet) && 
!holdsTransform(order.child))
+      outputPartitioning match {
+        case k: KeyedPartitioning if rest.nonEmpty =>
+          val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform))
+          prefix ++ rest.filter(order => keyExprs.contains(order.child))
+        case _ => prefix
+      }
+    }.getOrElse(Seq.empty)
+    outputPartitioning match {
+      case k: KeyedPartitioning
+          if reported.isEmpty && conf.v2BucketingPartitionKeyOrderingEnabled =>
         k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending))

Review Comment:
   Before this PR, a source that implements `SupportsReportOrdering` and 
returns an empty array (`Some(Nil)`) got no ordering, so Spark kept the sort 
above the scan. Now it gets `[key ASC]` by default, and so does a report that 
keeps no sort order.
   
   This is fine under the `KeyGroupedPartitioning` contract. But this suite's 
own `OrderAndPartitionAwareDataSource` breaks it: a split holds `i = [1, 1, 3]` 
under key `1`, and with `partitionKeys = j` a split holds `j = [6, 1, 2]` under 
key `2`. That is why `DataSourceV2Suite` had to turn the conf off. For such a 
source, a sort-merge join, a sort aggregate (`replaceHashWithSortAgg`), or 
`flatMapGroups` now reads unsorted input and can give wrong results, and the 
only workaround is the global conf.
   
   I see the PR description accepts this. Could the 4.4 migration note say it 
explicitly, for example that a source reporting an empty ordering is now also 
trusted to hold a single key per split? Then a connector owner who relied on an 
empty report to keep the sort knows to set 
`spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`.



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/DataSourceV2Suite.scala:
##########
@@ -315,7 +315,11 @@ class DataSourceV2Suite extends SharedSparkSession with 
AdaptiveSparkPlanHelper
   }
 
   test("ordering and partitioning reporting") {
-    withSQLConf(SQLConf.V2_BUCKETING_ENABLED.key -> "true") {
+    // The source breaks the `KeyGroupedPartitioning` contract, e.g. [1, 1, 3] 
under the key 1.
+    // So the ordering derived from the keys is turned off. This test checks 
the reported one.
+    withSQLConf(
+        SQLConf.V2_BUCKETING_ENABLED.key -> "true",
+        SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "false") {

Review Comment:
   Only the `(Some("i"), None)` case changes with this PR. The other cases, 
`[i]`, `[j]`, `[i,j]` for both the Scala and the Java source, now run only with 
the derivation off, so this suite no longer covers how a non-empty report 
interacts with the default conf.
   
   Could we scope the override to the affected case instead? For example, keep 
the outer `withSQLConf` as before and wrap only the `(Some("i"), None)` case in 
`withSQLConf(SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> 
"false")`. Making the source keep a single key per split would also work, but 
it changes the expected answers.



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -4078,6 +4078,41 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  test("SPARK-59981: a reported ordering that keeps no sort order falls back 
to the keys") {
+    withCustomReportingTable { reportingCatalog =>
+      reportingCatalog.reportedKeys = Seq(identity("id"))
+      // An empty report, one on `s`, which the join's projection prunes from 
the scan output, and
+      // one that starts with `s` and goes on with `data`, which is not a 
partition key.
+      val sOrder = sort(FieldReference("s"), SortDirection.ASCENDING)
+      val dataOrder = sort(FieldReference("data"), SortDirection.ASCENDING)
+      Seq(Seq.empty, Seq(sOrder), Seq(sOrder, dataOrder)).foreach { ordering =>

Review Comment:
   The three cases here all reduce to nothing, so the test pins the fallback 
but not its boundaries. Could we add a few cases?
   
   - `[data]`: `data` is in the output and not a key, so the report is kept as 
`[data]` and nothing is derived.
   - `[s, id]`: `id` survives past the pruned `s` through `rest.filter`, so the 
result is `[id]` from the report, not from the fallback. Its direction stays as 
reported, e.g. `[s, id DESC]` gives `[id DESC]`.
   - A `Some(Nil)` scan with two splits under one key, under a coalescing 
`GroupPartitionsExec`, keeps `[id]` through `preserveKeyOrderingOnCoalesce`.
   
   Without them, a change to the `reported.isEmpty` gate or to the 
`rest.filter` branch could change these orderings without failing any test.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExec.scala:
##########
@@ -377,8 +377,8 @@ case class GroupPartitionsExec(
   override def outputOrdering: Seq[SortOrder] = {
     if (!hasCoalescing) {
       // No coalescing: each output partition is exactly one input partition. 
The child's
-      // within-partition ordering is fully preserved (including any 
key-derived ordering that
-      // `DataSourceV2ScanExecBase` already prepended).
+      // within-partition ordering is fully preserved, including an ordering 
that
+      // `DataSourceV2ScanExecBase` derives from the partition keys.

Review Comment:
   `kWayMergeIsFeasible` (L258) is `child.outputOrdering.nonEmpty && 
childIsSafeForKWayMerge`. Before this PR only a `None` report made the derived 
key-only ordering reach it. Now `Some(Nil)` and fully dropped reports do too.
   
   With `spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled` on, 
`tryEnableSortedMerge` then plans a `SortedMergeCoalescedRDD` for such a scan. 
It merges on a key that is equal in every input of the group, so it gains no 
ordering beyond what `preserveKeyOrderingOnCoalesce` already keeps with plain 
concatenation. It still pays the priority-queue comparisons and turns off 
columnar output (`supportsColumnar` is false when `usesSortedMerge`).
   
   The conf is off by default, so this is minor. Should the merge be feasible 
only when the child's ordering has a sort order beyond the partition keys?



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -155,19 +161,21 @@ trait DataSourceV2ScanExecBase
    */
   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(order => order.references.subsetOf(outputSet) && 
!holdsTransform(order.child))
-        p match {
-          case k: KeyedPartitioning if rest.nonEmpty =>
-            val keyExprs = 
ExpressionSet(k.expressions.filterNot(holdsTransform))
-            prefix ++ rest.filter(order => keyExprs.contains(order.child))
-          case _ => prefix
-        }
-      case (_, k: KeyedPartitioning) if 
conf.v2BucketingPartitionKeyOrderingEnabled =>
+    val reported = ordering.map { o =>

Review Comment:
   Nit: spanning `Nil` already yields an empty prefix and rest, so `None` needs 
no separate path, and the two `outputPartitioning` matches can be one:
   
   ```scala
   val (prefix, rest) = ordering.getOrElse(Nil)
     .span(order => order.references.subsetOf(outputSet) && 
!holdsTransform(order.child))
   outputPartitioning match {
     case k: KeyedPartitioning =>
       val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform))
       val kept = prefix ++ rest.filter(order => keyExprs.contains(order.child))
       if (kept.isEmpty && conf.v2BucketingPartitionKeyOrderingEnabled) {
         k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending))
       } else {
         kept
       }
     case _ => prefix
   }
   ```
   
   This also makes it clear why `V2ScanPartitioningAndOrdering` no longer has 
to choose between `None` and `Some(Nil)`.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala:
##########
@@ -155,19 +161,21 @@ trait DataSourceV2ScanExecBase
    */
   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(order => order.references.subsetOf(outputSet) && 
!holdsTransform(order.child))
-        p match {
-          case k: KeyedPartitioning if rest.nonEmpty =>
-            val keyExprs = 
ExpressionSet(k.expressions.filterNot(holdsTransform))
-            prefix ++ rest.filter(order => keyExprs.contains(order.child))
-          case _ => prefix
-        }
-      case (_, k: KeyedPartitioning) if 
conf.v2BucketingPartitionKeyOrderingEnabled =>
+    val reported = ordering.map { o =>
+      val (prefix, rest) =
+        o.span(order => order.references.subsetOf(outputSet) && 
!holdsTransform(order.child))
+      outputPartitioning match {
+        case k: KeyedPartitioning if rest.nonEmpty =>
+          val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform))
+          prefix ++ rest.filter(order => keyExprs.contains(order.child))
+        case _ => prefix
+      }
+    }.getOrElse(Seq.empty)
+    outputPartitioning match {
+      case k: KeyedPartitioning
+          if reported.isEmpty && conf.v2BucketingPartitionKeyOrderingEnabled =>

Review Comment:
   Not blocking, but I think this gate makes the plan depend on column pruning. 
With key `identity(id)` and a report `[s]`:
   
   - `SELECT t1.data, t2.data ... JOIN ON t1.id = t2.id` prunes `s`, so the 
report is dropped and the scan derives `[id ASC]`. No sort.
   - `SELECT t1.s, t2.data ...` keeps `s`, so the scan reports `[s]` and the 
join adds a sort on `id` again.
   - A scan merged by `PlanMerger` whose widened output keeps `s` loses `[id 
ASC]` too. `mergeDegradesReporting` compares only the reported `ordering` 
(`[s]` on both sides), so it does not see this.
   
   A report `[data]` or `[id DESC]` behaves the same way: any surviving sort 
order suppresses the key ordering, even though `id` is constant in each split.
   
   Prepending the key orders (`keyOrders ++ reported`, deduped), as the old 
`GroupPartitionsExec` comment described ("key-derived ordering ... already 
prepended"), would fix these cases. But since `orderingSatisfies` is 
prefix-based, it would stop `[data]` from satisfying a required `[data]`. So 
the root fix is probably to make the satisfaction check aware that the key 
expressions are constant within a split. That seems out of scope here. Should 
we file a separate JIRA for it, and maybe mention the pruning dependency in the 
Scaladoc?



-- 
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]

Reply via email to