andygrove commented on code in PR #5331:
URL: https://github.com/apache/datafusion-comet/pull/5331#discussion_r4073739836


##########
docs/source/user-guide/latest/iceberg.md:
##########
@@ -180,6 +193,12 @@ The following scenarios will fall back to the JVM Iceberg 
reader:
 - Scans with residual filters using `truncate`, `bucket`, `year`, `month`, 
`day`, or `hour`
   transform functions (partition pruning still works, but row-level filtering 
of these
   transforms falls back)
+- Sorted tables where the sort key is a transform (e.g. `bucket`, `truncate`) 
rather than a plain
+  column: Comet reports only identity sort keys today, so the scan falls back 
to Spark to preserve
+  the reported ordering
+- Sorted tables whose sort key is a `UUID` column: Iceberg sorts UUID by its 
own byte comparator

Review Comment:
   The transform sort key and the UUID sort key landing here is exactly what I 
was after last round, thank you.
   
   The unprojected sort key belongs in this list too, since as things stand 
today that is also a fallback to Spark and it will be the most common one 
people hit. Whatever we settle on in the `CometScanRule` comment, the list 
should describe it.



##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala:
##########
@@ -893,6 +894,82 @@ object CometIcebergNativeScan extends 
CometOperatorSerde[CometBatchScanExec] wit
     Some(builder.setIcebergScan(icebergScanBuilder).build())
   }
 
+  /**
+   * The part of an Iceberg-reported sort order that the native per-partition 
merge can honour, or
+   * Nil when the merge must stay off. Two callers use this one gate: the 
proto serialization
+   * (which turns on the native SortPreservingMergeExec) and
+   * CometIcebergNativeScanExec.outputOrdering (which tells Spark the scan is 
sorted). Sharing the
+   * gate means the two always agree.
+   *
+   * v1 accepts only identity sort fields on top-level columns that are in the 
projection. Each
+   * SortOrder child must be an AttributeReference in `output`, and must 
serialize to proto.
+   * Transform sort fields (bucket/truncate/...) are not AttributeReferences, 
so they fall through
+   * to Nil and we read unordered. Checking exprToProto here, not just in the 
proto path, keeps
+   * the two callers in step: outputOrdering never advertises an order the 
proto path would drop.
+   *
+   * We trust Iceberg on file-level sortedness. If it reports an ordering, 
SortOrderAnalyzer has
+   * already checked each file's sort_order_id matches the table order, so 
every file is sorted.
+   *
+   * We read scanExec.ordering (the raw reported order), not 
scanExec.outputOrdering. Spark blanks

Review Comment:
   This docstring explains the choice of `scanExec.ordering` over 
`scanExec.outputOrdering` with "Spark blanks outputOrdering when a partition 
holds more than one file", and I think that has the mechanism backwards in a 
way worth fixing, because this is the safety-critical read.
   
   Spark blanks it when multiple `InputPartition`s share a partition key, not 
when a partition holds multiple files. 
`DataSourceV2ScanExecBase.outputOrdering` is `ordering.filter(_ => 
groupedPartitions.forall(_.groupedParts.forall(_.parts.length <= 1)))`. 
Iceberg's `SortOrderAnalyzer.hasUniquePartitionKeys` refuses to report in 
exactly that same case, with a comment giving the same reason. So on the 
Iceberg path the two values are always equal, and the 
many-files-in-one-partition case this merge targets is one where Spark does not 
blank at all.
   
   That means reading `outputOrdering` instead would be behaviour-preserving 
today, and it would stop us depending on an Iceberg-side invariant to stay 
safe. If a future Iceberg relaxes `hasUniquePartitionKeys`, or another source 
reports an ordering Spark then blanks, reading `ordering` has us pay a k-way 
merge, or a full spillable `SortExec` above the cap, underneath a `Sort` that 
Spark kept anyway. Would you rather switch the read, or keep `ordering` and 
have the comment name the Iceberg check it is leaning on?



##########
spark/src/test/scala/org/apache/comet/CometIcebergSortMergeReadSuite.scala:
##########
@@ -0,0 +1,885 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet
+
+import java.util.concurrent.atomic.AtomicInteger
+
+import org.apache.spark.sql.CometTestBase
+import org.apache.spark.sql.catalyst.expressions.{Add, Ascending, 
AttributeReference, Literal, SortOrder}
+import org.apache.spark.sql.comet.{CometIcebergNativeScanExec, CometSortExec}
+import org.apache.spark.sql.comet.execution.shuffle.CometShuffleExchangeExec
+import org.apache.spark.sql.execution.{SortExec, SparkPlan}
+import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper
+import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
+import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec
+import org.apache.spark.sql.types.IntegerType
+
+import org.apache.comet.CometSparkSessionExtensions.isSpark40Plus
+import org.apache.comet.serde.operator.CometIcebergNativeScan
+
+/**
+ * Tests for the sort-aware native Iceberg scan (branch `stream-merge`): the 
scan reports the
+ * Iceberg table sort order to Spark and does a per-partition streaming k-way 
merge of the
+ * already-sorted files, so Catalyst can drop the Sort and (with 
storage-partitioned join) the
+ * Exchange that a sort-merge join / grouped aggregate / window would 
otherwise need.
+ *
+ * The suite has two parts:
+ *   - Unit tests of the [[CometIcebergNativeScan.reportableOrdering]] gate 
(no SparkSession
+ *     required) -- the single decision shared by the proto serialization 
(which turns on the
+ *     native SortPreservingMergeExec) and 
CometIcebergNativeScanExec.outputOrdering (which tells
+ *     Spark the scan is sorted).
+ *   - End-to-end tests over real Iceberg tables that exercise every Spark 4.0 
mechanism which
+ *     exploits already-sorted input to avoid a Sort or a shuffle (all keyed 
off
+ *     `SortOrder.orderingSatisfies`, prefix semantics): 
`SupportsReportOrdering` ->
+ *     `BatchScanExec.outputOrdering` and `EnsureRequirements` eliding the 
required-child Sort;
+ *     `EliminateSorts` / `RemoveRedundantSorts`; 
`SortMergeJoinExec.requiredChildOrdering` +
+ *     storage-partitioned join (`KeyGroupedPartitioning`,
+ *     `spark.sql.sources.v2.bucketing.enabled`); `ReplaceHashWithSortAgg` / 
`SortAggregateExec`;
+ *     `WindowExec`; `TakeOrderedAndProjectExec`.
+ *
+ * Two invariants determine the end-to-end assertions:
+ *   1. Correctness is checked unconditionally via `checkSparkAnswer` (Comet 
vs vanilla Spark).
+ *      This is the primary guarantee: any k-way-merge defect (dropped, 
duplicated or mis-ordered
+ *      rows, or an outputOrdering/outputPartitioning that does not match the 
rows the native
+ *      operator actually produces) shows up as a result mismatch. It holds on 
any Iceberg build.
+ *      2. The strict "no Sort / no Exchange" plan assertions are the *target* 
of this feature.
+ *      They only hold where the Iceberg build actually reports the ordering
+ *      (`SupportsReportOrdering`, today an Iceberg fork feature -- the 
published/upstream Iceberg
+ *      used in CI does not report it) and where the native scan reports 
`KeyGroupedPartitioning`.
+ *      Each such test therefore runs the correctness check first, then 
`assume`s the reporting is
+ *      active before asserting the plan shape, so it enforces the contract on 
a reporting build
+ *      and is skipped (not failed) elsewhere. `sort = 0` is asserted only for 
*operator-required*
+ *      orderings (SMJ / aggregate / window), never for a global `ORDER BY`, 
which Spark keeps
+ *      regardless (the per-partition merge is not a global order).
+ *
+ * These cannot be SQL-file fixtures: setting an Iceberg sort order needs the 
Iceberg Java API
+ * (the Comet test session registers no Iceberg SQL extensions, so `WRITE 
ORDERED BY` will not
+ * parse), and the plan-shape assertions need access to the executed SparkPlan.
+ *
+ * Each test gets its own catalog name and temp warehouse (the Hadoop 
`SparkCatalog` instance is
+ * cached per catalog name, so a shared name would bind every test to the 
first warehouse), and
+ * its tables are dropped in a `finally` so a failing test cannot leak a table 
into a later one.
+ */
+class CometIcebergSortMergeReadSuite
+    extends CometTestBase
+    with CometIcebergTestBase
+    with AdaptiveSparkPlanHelper {
+
+  // 
---------------------------------------------------------------------------------------------
+  // Unit tests of the reportableOrdering gate (merged from the former 
CometIcebergSortMergeSuite).
+  // v1 identity-scope gate as a pure function, independent of a live 
SparkSession or an Iceberg
+  // build that reports ordering. The flag defaults to enabled, so no SQLConf 
override is needed.
+  // 
---------------------------------------------------------------------------------------------
+
+  private val gateA = AttributeReference("a", IntegerType)()
+  private val gateB = AttributeReference("b", IntegerType)()
+
+  test("gate: identity ordering on projected columns is reportable") {
+    val ordering = Seq(SortOrder(gateA, Ascending))
+    assert(
+      CometIcebergNativeScan.reportableOrdering(Some(ordering), Seq(gateA, 
gateB)) === ordering)
+  }
+
+  test("gate: ordering on a column outside the projection falls back") {
+    val ordering = Seq(SortOrder(gateA, Ascending))
+    assert(CometIcebergNativeScan.reportableOrdering(Some(ordering), 
Seq(gateB)).isEmpty)
+  }
+
+  test("gate: a transform (non-AttributeReference) sort child falls back") {
+    val ordering = Seq(SortOrder(Add(gateA, Literal(1)), Ascending))
+    assert(CometIcebergNativeScan.reportableOrdering(Some(ordering), 
Seq(gateA)).isEmpty)
+  }
+
+  test("gate: if any sort field is unreportable, the whole ordering falls 
back") {
+    val ordering = Seq(SortOrder(gateA, Ascending), SortOrder(Add(gateB, 
Literal(1)), Ascending))
+    assert(CometIcebergNativeScan.reportableOrdering(Some(ordering), 
Seq(gateA, gateB)).isEmpty)
+  }
+
+  test("gate: absent or empty ordering falls back") {
+    assert(CometIcebergNativeScan.reportableOrdering(None, Seq(gateA)).isEmpty)
+    assert(CometIcebergNativeScan.reportableOrdering(Some(Seq.empty), 
Seq(gateA)).isEmpty)
+  }
+
+  test("gate: a sort key on an ordering-unsafe column (e.g. UUID) falls back") 
{
+    // "a" stands in for a UUID column: Iceberg maps it to StringType but 
sorts by its own order,
+    // so it is unsafe to honour even though it looks like a plain string at 
the Spark level.
+    val ordering = Seq(SortOrder(gateA, Ascending))
+    assert(
+      CometIcebergNativeScan
+        .reportableOrdering(Some(ordering), Seq(gateA, gateB), Set("a"))
+        .isEmpty)
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // End-to-end fixtures.
+  // 
---------------------------------------------------------------------------------------------
+
+  private val catalogCounter = new AtomicInteger(0)
+
+  // Storage-partitioned join config. preserve-data-grouping + v2 bucketing 
let Iceberg report
+  // KeyGroupedPartitioning on the BatchScanExec, and the join knobs force a 
sort-merge join over
+  // co-partitioned inputs so the Exchange can be eliminated. Comet does not 
report partitioning
+  // itself; Spark's EnsureRequirements eliminates the shuffle on the 
BatchScanExec before Comet
+  // converts the scan. Adaptive off keeps the executed plan stable for the 
counts below.
+  private val spjConf: Seq[(String, String)] = Seq(
+    "spark.sql.iceberg.planning.preserve-data-ordering" -> "true",
+    "spark.sql.iceberg.planning.preserve-data-grouping" -> "true",
+    "spark.sql.sources.v2.bucketing.enabled" -> "true",
+    "spark.sql.sources.v2.bucketing.pushPartValues.enabled" -> "true",
+    "spark.sql.requireAllClusterKeysForCoPartition" -> "false",
+    "spark.sql.autoBroadcastJoinThreshold" -> "-1",
+    "spark.sql.join.preferSortMergeJoin" -> "true",
+    "spark.sql.adaptive.enabled" -> "false")
+
+  /**
+   * Runs `f` against a fresh, uniquely-named Hadoop catalog backed by a fresh 
temp warehouse,
+   * then drops the named tables (IF EXISTS) in a `finally` -- so tables are 
cleaned up even when
+   * the test body fails, and no two tests can collide on a table name. `f` 
receives the catalog
+   * name; tables live under the `db` namespace, e.g. `$cat.db.$table`.
+   */
+  private def withSortedTables(extraConf: Seq[(String, String)])(tables: 
String*)(
+      f: String => Unit): Unit = {
+    assume(icebergAvailable, "Iceberg not available in classpath")
+    withTempIcebergDir { warehouseDir =>
+      val cat = s"sort_cat_${catalogCounter.incrementAndGet()}"
+      val cometConf = Seq(
+        s"spark.sql.catalog.$cat" -> "org.apache.iceberg.spark.SparkCatalog",
+        s"spark.sql.catalog.$cat.type" -> "hadoop",
+        s"spark.sql.catalog.$cat.warehouse" -> warehouseDir.getAbsolutePath,
+        CometConf.COMET_ENABLED.key -> "true",
+        CometConf.COMET_EXEC_ENABLED.key -> "true",
+        CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true")
+      withSQLConf((cometConf ++ extraConf): _*) {
+        try f(cat)
+        finally tables.foreach(t => spark.sql(s"DROP TABLE IF EXISTS 
$cat.db.$t"))
+      }
+    }
+  }
+
+  /**
+   * Sets the table sort order via the Iceberg Java API. `cols` is (column, 
ascending); ascending
+   * uses Iceberg's default NULLS FIRST and descending its default NULLS LAST, 
matching the ORDER
+   * BY null-ordering the tests use.
+   */
+  private def replaceSortOrder(
+      cat: String,
+      namespace: String,
+      table: String,
+      cols: (String, Boolean)*): Unit = {
+    val catalog = spark.sessionState.catalogManager
+      .catalog(cat)
+      .asInstanceOf[org.apache.iceberg.spark.SparkCatalog]
+    val ident =
+      org.apache.spark.sql.connector.catalog.Identifier.of(Array(namespace), 
table)
+    val icebergTable = catalog
+      .loadTable(ident)
+      .asInstanceOf[org.apache.iceberg.spark.source.SparkTable]
+      .table()
+    var sortOrder = icebergTable.replaceSortOrder()
+    cols.foreach { case (c, asc) =>
+      sortOrder = if (asc) sortOrder.asc(c) else sortOrder.desc(c)
+    }
+    sortOrder.commit()
+  }
+
+  /**
+   * Sets a transform (bucket) sort order via the Iceberg Java API. Iceberg 
orders files by the
+   * bucket hash, so the reported sort field is a transform expression, not a 
plain column. Spark
+   * 4.0+ can convert a bucket transform ordering 
(V2ScanPartitioningAndOrdering threads the
+   * function catalog and V2ExpressionUtils special-cases BucketTransform); 
Spark 3.4 cannot, so
+   * the caller gates the test on isSpark40Plus. Either way Comet's 
reportableOrdering rejects the
+   * non-AttributeReference sort child (v1 is identity-scope only, #5339), so 
Comet must fall
+   * back.
+   */
+  private def replaceSortOrderBucket(
+      cat: String,
+      namespace: String,
+      table: String,
+      col: String,
+      numBuckets: Int): Unit = {
+    val catalog = spark.sessionState.catalogManager
+      .catalog(cat)
+      .asInstanceOf[org.apache.iceberg.spark.SparkCatalog]
+    val ident =
+      org.apache.spark.sql.connector.catalog.Identifier.of(Array(namespace), 
table)
+    val icebergTable = catalog
+      .loadTable(ident)
+      .asInstanceOf[org.apache.iceberg.spark.source.SparkTable]
+      .table()
+    icebergTable
+      .replaceSortOrder()
+      .asc(org.apache.iceberg.expressions.Expressions.bucket(col, numBuckets))
+      .commit()
+  }
+
+  /**
+   * Each string becomes a separate INSERT, hence a separate data file, so 
merging is required.
+   */
+  private def insertBatches(cat: String, table: String, batches: String*): 
Unit =
+    batches.foreach(values => spark.sql(s"INSERT INTO $cat.db.$table VALUES 
$values"))
+
+  private def nativeScans(plan: SparkPlan): Seq[CometIcebergNativeScanExec] =
+    collect(stripAQEPlan(plan)) { case s: CometIcebergNativeScanExec => s }
+
+  private def countSorts(plan: SparkPlan): Int =
+    collect(stripAQEPlan(plan)) {
+      case s: SortExec => s
+      case s: CometSortExec => s
+    }.size
+
+  private def countShuffles(plan: SparkPlan): Int =
+    collect(stripAQEPlan(plan)) {
+      case e: ShuffleExchangeExec => e
+      case e: CometShuffleExchangeExec => e
+    }.size
+
+  /** True once every native scan in the plan advertises the reported 
ordering. */
+  private def orderingReported(plan: SparkPlan): Boolean = {
+    val scans = nativeScans(plan)
+    scans.nonEmpty && scans.forall(_.outputOrdering.nonEmpty)
+  }
+
+  // NOTE ON CANCELED TESTS: the helper below CANCELS the test (ScalaTest 
`assume`, reported as
+  // "!!! CANCELED !!!", not a failure) when the Iceberg build on the 
classpath does not implement
+  // the DSv2 `SupportsReportOrdering` API. Published/upstream Iceberg (the 
runtime used in CI and
+  // the default mvn profiles) does not report a sort order, so 
`outputOrdering` comes back empty
+  // and there is no eliminated Sort to assert on. These tests therefore show 
as canceled there --
+  // that is expected, NOT a regression. The preceding `checkSparkAnswer` has 
already validated
+  // correctness; only the sort/shuffle-elimination plan assertion is skipped. 
Run against an
+  // ordering-reporting (fork) Iceberg build to exercise those assertions.
+
+  /** Cancels the test (see note above) unless the scan reported an ordering. 
*/
+  private def assumeOrderingReported(plan: SparkPlan): Unit =
+    assume(
+      orderingReported(plan),
+      "current Iceberg build does not implement SupportsReportOrdering (no 
ordering reported); " +
+        "sort-elimination assertion skipped")
+
+  /**
+   * True if the Iceberg build on the classpath reports a sort order for 
`query`. Determined
+   * independently of Comet: the query is planned with the native Iceberg scan 
disabled, and we
+   * check whether Spark's own BatchScanExec advertises an outputOrdering. 
This is the signal that
+   * makes the fallback checks below meaningful -- when Iceberg does not 
report (the published
+   * Iceberg in CI), there is no reported ordering for Comet to decline, so 
those checks are
+   * skipped rather than firing on a scan that was never eligible to fall back.
+   */
+  private def icebergReportsOrdering(query: String): Boolean = {
+    // withSQLConf returns the block value on Spark 4.x but Unit on 3.4/3.5, 
so capture in a var
+    // rather than relying on its return value.
+    var reported = false
+    withSQLConf(CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "false") {
+      val plan = spark.sql(query).queryExecution.executedPlan
+      reported = collect(stripAQEPlan(plan)) { case b: BatchScanExec => b }
+        .exists(_.outputOrdering.nonEmpty)
+    }
+    reported
+  }
+
+  /**
+   * The Spark-fallback contract for a reported-but-unhonorable ordering: when 
the native scan
+   * cannot guarantee an ordering Iceberg reported, Comet must not convert the 
scan at all.
+   * Reading unordered natively would return silently-wrong results, because 
EnsureRequirements
+   * has already dropped the Sort above the scan on the strength of Iceberg's 
report. So the plan
+   * must contain no CometIcebergNativeScanExec -- the scan stays on Spark's 
Iceberg reader. Gated
+   * on a reporting Iceberg build (see icebergReportsOrdering); correctness is 
asserted separately
+   * by the caller.
+   */
+  private def assertFellBackToSpark(query: String, plan: SparkPlan): Unit = {
+    assume(
+      icebergReportsOrdering(query),
+      "current Iceberg build does not report an ordering; the Spark-fallback 
path is not exercised")
+    assert(
+      nativeScans(plan).isEmpty,
+      s"Comet reported an ordering it cannot honour instead of falling back to 
Spark:\n$plan")
+  }
+
+  // The reporting mechanism: SupportsReportOrdering -> 
CometIcebergNativeScanExec.outputOrdering
+  //
+  // NOTE ON TABLE SHAPE: Iceberg only reports a sort order for a table with a 
non-empty grouping
+  // key -- isOrderingEnabled in apache/iceberg#16750 is 
`!groupingKeyType().fields().isEmpty() &&
+  // canReportOrdering(...)`, and Iceberg's own 
testNoMergeReaderForUnpartitionedSortedTable asserts
+  // an unpartitioned sorted table reports nothing. So a merge test MUST use a 
partitioned table, or
+  // Iceberg reports no ordering, Comet takes the plain unordered read, and 
the merge never runs (the
+  // checkSparkAnswer then passes by comparing the unordered path against 
itself). Every merge test
+  // below partitions by a single-value column `p` so the ordering is reported 
while all files stay
+  // in one partition -- the shape that actually exercises the multi-file 
k-way merge -- and asserts
+  // assumeOrderingReported after the correctness check so a reporting build 
proves the merge ran.
+
+  test("native scan reports the table sort order for a multi-file sorted 
table") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      insertBatches(cat, "t", "(1,'a','P1'),(3,'c','P1')", 
"(2,'b','P1'),(4,'d','P1')")
+
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id")
+      assume(nativeScans(plan).nonEmpty, "query did not use the native Iceberg 
scan")
+      assumeOrderingReported(plan)
+      assert(
+        
nativeScans(plan).head.outputOrdering.head.child.references.exists(_.name == 
"id"),
+        s"expected the scan to report an ordering on id:\n$plan")
+    }
+  }
+
+  test("sort-merge disabled keeps the scan native and still reports the 
ordering (via sort)") {
+    withSortedTables(spjConf ++ 
Seq(CometConf.COMET_ICEBERG_SORT_MERGE_ENABLED.key -> "false"))(
+      "t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      insertBatches(cat, "t", "(1,'a','P1'),(3,'c','P1')", 
"(2,'b','P1'),(4,'d','P1')")
+
+      // Disabling only turns off the k-way merge, not the whole native scan: 
Comet still reads
+      // the table and still honours the reported order (via a spillable 
sort). Correctness must
+      // hold, and on an ordering-reporting Iceberg build the scan still 
advertises the order.
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id")
+      assert(
+        nativeScans(plan).nonEmpty,
+        s"sort-merge disabled must not force the scan back to Spark:\n$plan")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  // K-way merge correctness. The tables are partitioned so Iceberg reports 
the ordering and the
+  // merge actually runs (see the NOTE ON TABLE SHAPE above); 
assumeOrderingReported proves that on a
+  // reporting build. Order sensitivity is verified with a window over the 
merged input rather than a
+  // global ORDER BY, because Spark keeps its final Sort for a global ORDER BY 
(a per-partition order
+  // does not satisfy a global one) and that re-sort would repair -- and hide 
-- a mis-ordered merge.
+
+  test("merges multiple sorted files per partition and preserves order") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (c1 INT, c2 INT, data STRING) USING iceberg " 
+
+          "PARTITIONED BY (bucket(4, c1))")
+      replaceSortOrder(cat, "db", "t", "c1" -> true, "c2" -> true)
+      // Multiple files, each holding rows for every c1 bucket, so the k-way 
merge runs within a
+      // partition; c2 interleaves across the files so a mis-ordered merge 
changes the row numbering.
+      insertBatches(cat, "t", "(1,1,'a'),(2,1,'b')", "(1,3,'c'),(2,3,'d')", 
"(1,2,'e'),(2,2,'f')")
+
+      // ROW_NUMBER over the reported (c1, c2) order: Spark keeps that order 
(the window's required
+      // ordering is satisfied, no re-sort), so a mis-ordered merge yields 
wrong row numbers.
+      val query =
+        s"SELECT c1, c2, ROW_NUMBER() OVER (PARTITION BY c1 ORDER BY c2) AS rn 
FROM $cat.db.t"
+      val (_, plan) = checkSparkAnswer(query)
+      assume(nativeScans(plan).nonEmpty, "query did not use the native Iceberg 
scan")
+      assumeOrderingReported(plan)
+      assert(
+        countSorts(plan) == 0,
+        s"the merge must satisfy the window ordering without a 
re-sort:\n$plan")
+    }
+  }
+
+  test("merge interleaves duplicate sort-key values across files") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      // The same id appears in several files; the merge must keep every row, 
not drop or mis-order.
+      insertBatches(
+        cat,
+        "t",
+        "(1,'a','P1'),(2,'b','P1')",
+        "(1,'c','P1'),(2,'d','P1')",
+        "(1,'e','P1'),(3,'f','P1')")
+
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id, data")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  test("merge applies merge-on-read deletes across a multi-file sorted 
partition") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
" +
+          "PARTITIONED BY (p) " +
+          "TBLPROPERTIES ('format-version'='2', 
'write.delete.mode'='merge-on-read')")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      // Several files so a merge is required; then delete rows from some of 
them. On a v2
+      // merge-on-read table DELETE writes delete files rather than rewriting 
the data files, so
+      // the scan must apply the deletes while merging the still-sorted files.
+      insertBatches(
+        cat,
+        "t",
+        "(1,'a','P1'),(4,'d','P1')",
+        "(2,'b','P1'),(5,'e','P1')",
+        "(3,'c','P1'),(6,'f','P1')")
+      spark.sql(s"DELETE FROM $cat.db.t WHERE id IN (2, 5)")
+
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  test("merge honours NULLS FIRST on an ascending sort key") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg 
" +
+          "PARTITIONED BY (c3)")
+      replaceSortOrder(cat, "db", "t", "c1" -> true) // ASC -> Iceberg default 
NULLS FIRST
+      insertBatches(
+        cat,
+        "t",
+        "(null,'x','P1'),(3,'c','P1')",
+        "(null,'y','P1'),(1,'a','P1'),(2,'b','P1')")
+
+      checkSparkAnswer(
+        s"SELECT c1, c2 FROM $cat.db.t WHERE c3 = 'P1' ORDER BY c1 ASC NULLS 
FIRST, c2")
+    }
+  }
+
+  test("merge honours NULLS LAST on a descending sort key") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg 
" +
+          "PARTITIONED BY (c3)")
+      replaceSortOrder(cat, "db", "t", "c1" -> false) // DESC -> Iceberg 
default NULLS LAST
+      insertBatches(
+        cat,
+        "t",
+        "(null,'x','P1'),(1,'a','P1')",
+        "(null,'y','P1'),(3,'c','P1'),(2,'b','P1')")
+
+      checkSparkAnswer(
+        s"SELECT c1, c2 FROM $cat.db.t WHERE c3 = 'P1' ORDER BY c1 DESC NULLS 
LAST, c2")
+    }
+  }
+
+  test("merge on a descending sort order preserves order") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (c1 INT, c2 INT, data STRING) USING iceberg " 
+
+          "PARTITIONED BY (bucket(4, c1))")
+      replaceSortOrder(cat, "db", "t", "c1" -> true, "c2" -> false) // c2 
DESC, Iceberg NULLS LAST
+      insertBatches(cat, "t", "(1,3,'a'),(2,3,'b')", "(1,1,'c'),(2,1,'d')", 
"(1,2,'e'),(2,2,'f')")
+
+      // Window ordered DESC over the merged (c1, c2 DESC) order; a 
mis-ordered descending merge
+      // changes the row numbers (a global ORDER BY DESC would be re-sorted by 
Spark and hide it).
+      val query =
+        s"SELECT c1, c2, ROW_NUMBER() OVER (PARTITION BY c1 ORDER BY c2 DESC) 
AS rn FROM $cat.db.t"
+      val (_, plan) = checkSparkAnswer(query)
+      assume(nativeScans(plan).nonEmpty, "query did not use the native Iceberg 
scan")
+      assumeOrderingReported(plan)
+      assert(
+        countSorts(plan) == 0,
+        s"the descending merge must satisfy the window ordering without a 
re-sort:\n$plan")
+    }
+  }
+
+  test("merge on a multi-column sort order") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING, p STRING) 
USING iceberg " +
+          "PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "c3" -> true, "c1" -> true)
+      insertBatches(
+        cat,
+        "t",
+        "(1,'a','A','P1'),(3,'c','A','P1')",
+        "(2,'b','A','P1'),(1,'a','B','P1')",
+        "(2,'b','B','P1'),(3,'c','B','P1')")
+
+      val (_, plan) = checkSparkAnswer(s"SELECT c3, c1, c2 FROM $cat.db.t 
ORDER BY c3, c1")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  test("single file needs no merge") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      insertBatches(cat, "t", "(1,'a','P1'),(2,'b','P1')")
+
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  test("many small files in one partition merge correctly") {
+    // Raise the per-partition file limit so the k-way merge (not the sort 
fallback) runs with many
+    // files -- the shape this feature targets, where the merge opens one 
reader per file at once.
+    // Partitioned by a single value so Iceberg reports the ordering and the 
~70-way merge actually
+    // runs; assumeOrderingReported proves it on a reporting build. 
checkSparkAnswer guards
+    // correctness; asserting on memory-pool usage / peak concurrent readers 
is a TODO for #5343.
+    val conf =
+      spjConf :+ 
(CometConf.COMET_ICEBERG_SORT_MERGE_MAX_FILES_PER_PARTITION.key -> "1000")
+    withSortedTables(conf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      val rows = (1 to 70).map(i => s"($i,'v$i','P1')")
+      insertBatches(cat, "t", rows: _*)
+
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  test("many small files in one partition stay correct (exercises the sort 
fallback)") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (id INT, data STRING, p STRING) USING iceberg 
PARTITIONED BY (p)")
+      replaceSortOrder(cat, "db", "t", "id" -> true)
+      // 70 single-row inserts -> 70 files in one partition, above the default 
maxFilesPerPartition
+      // (64). Partitioned so Iceberg reports the ordering; above the cap the 
native scan takes the
+      // fallback -- a single unordered read plus a spillable SortExec, not a 
70-way merge -- which
+      // is the shape this feature targets (a sorted table with many small 
commits). assumeOrdering
+      // Reported proves the ordering was reported (so the sort-fallback path 
really ran) on a
+      // reporting build. checkSparkAnswer guards correctness; asserting on 
memory-pool usage / peak
+      // concurrent readers is a TODO for #5343.
+      val rows = (1 to 70).map(i => s"($i,'v$i','P1')")
+      insertBatches(cat, "t", rows: _*)
+
+      val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER 
BY id")
+      assumeOrderingReported(plan)
+    }
+  }
+
+  test("partitioned table with several files per partition") {
+    withSortedTables(spjConf)("t") { cat =>
+      spark.sql(
+        s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg 
" +
+          "PARTITIONED BY (c3)")
+      replaceSortOrder(cat, "db", "t", "c1" -> true)
+      insertBatches(
+        cat,
+        "t",
+        "(1,'a','P1'),(3,'c','P1')",
+        "(2,'b','P1'),(4,'d','P1')",
+        "(5,'e','P2'),(7,'g','P2')",
+        "(6,'f','P2'),(8,'h','P2')")
+
+      checkSparkAnswer(s"SELECT c1, c2 FROM $cat.db.t WHERE c3 = 'P1' ORDER BY 
c1")
+      checkSparkAnswer(s"SELECT c1, c2 FROM $cat.db.t WHERE c3 = 'P2' ORDER BY 
c1")
+    }
+  }
+
+  test("sort key absent from the projection: falls back to an unordered read 
but stays correct") {

Review Comment:
   Following on from the comment on `CometScanRule`, I do not think this test 
pins the behaviour it is named for.
   
   ```scala
   val (_, plan) = checkSparkAnswer(s"SELECT c2 FROM $cat.db.t WHERE c3 = 'P1'")
   nativeScans(plan).foreach { scan =>
     assert(scan.outputOrdering.isEmpty, ...)
   }
   ```
   
   On a reporting build this query falls back to Spark, so `nativeScans(plan)` 
is empty and the `foreach` asserts nothing. On the published Iceberg nothing 
was reported in the first place, so the assertion is trivially true there too. 
Either way it passes, and the title says "falls back to an unordered read" 
while the behaviour is a fallback to Spark.
   
   Could this assert the outcome directly, whichever one we settle on above? If 
the scan should stay native with no ordering, 
`assert(nativeScans(plan).nonEmpty)` first and then the 
`outputOrdering.isEmpty` check. If it should stay a Spark fallback, 
`assertFellBackToSpark` already exists and says so.



##########
native/core/src/execution/operators/iceberg_scan.rs:
##########
@@ -175,18 +239,14 @@ impl IcebergScanExec {
     fn execute_with_tasks(
         &self,
         tasks: Vec<FileScanTask>,
+        partition: usize,
         context: Arc<TaskContext>,
     ) -> DFResult<SendableRecordBatchStream> {
         let output_schema = Arc::clone(&self.output_schema);
-        let file_io = load_file_io(
-            &self.catalog_properties,
-            &self.metadata_location,
-            &self.catalog_name,
-            AccessMode::Read,
-        )?;
+        let file_io = self.file_io.clone();

Review Comment:
   On the per-file `CachingDeleteFileLoader`, you said you would fold this into 
#5343. Reading that issue, its body is entirely about bounding how many files 
the merge opens at once, and it does not mention the delete loader being 
rebuilt per `execute`, or the measurement: `bytes_scanned` going from 2717 to 
4256 with two data files sharing one positional delete file, and up to 64x on a 
merge-on-read table at the default cap.
   
   Could you add it to #5343, or file it separately? Once this lands, the shape 
of the regression is not recoverable from the code without redoing the 
measurement.



##########
native/core/src/execution/operators/iceberg_scan.rs:
##########
@@ -707,6 +768,141 @@ mod tests {
             .unwrap();
     }
 
+    fn int_schema() -> arrow::datatypes::SchemaRef {
+        use arrow::datatypes::{DataType, Field, Schema};
+        Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]))
+    }
+
+    fn single_col_ordering() -> Option<datafusion::physical_expr::LexOrdering> 
{
+        use arrow::compute::SortOptions;
+        use datafusion::physical_expr::expressions::Column;
+        use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr};
+        LexOrdering::new(vec![PhysicalSortExpr {
+            expr: Arc::new(Column::new("a", 0)),
+            options: SortOptions::default(),
+        }])
+    }
+
+    // Builds a scan over three empty-delete tasks with the given reported 
ordering.
+    fn exec_with_ordering(
+        ordering: Option<datafusion::physical_expr::LexOrdering>,
+    ) -> IcebergScanExec {
+        use std::collections::HashMap;
+        let tasks = vec![
+            task_with_deletes(vec![]),
+            task_with_deletes(vec![]),
+            task_with_deletes(vec![]),
+        ];
+        IcebergScanExec::new(
+            "metadata.json".to_string(),
+            int_schema(),
+            HashMap::new(),
+            "cat".to_string(),
+            tasks,
+            1,
+            ordering,
+        )
+        .unwrap()
+    }
+
+    // A reported ordering turns the scan into a multi-partition operator (one 
partition per task)
+    // so a SortPreservingMergeExec above can k-way merge the per-file sorted 
streams.
+    #[test]
+    fn reported_ordering_makes_scan_multi_partition() {
+        let exec = exec_with_ordering(single_col_ordering());
+        assert_eq!(exec.properties().partitioning.partition_count(), 3);
+    }
+
+    // Without a reported ordering the scan stays single-partition (Comet 
drives only execute(0),
+    // which must read every task), preserving the legacy unordered behaviour.
+    #[test]
+    fn no_ordering_keeps_single_partition() {
+        let exec = exec_with_ordering(None);
+        assert_eq!(exec.properties().partitioning.partition_count(), 1);
+    }
+
+    // The ordered scan reads each file as its own sorted partition and relies 
on
+    // SortPreservingMergeExec to k-way merge them into one globally sorted 
stream. This feeds known
+    // sorted partitions (with duplicate keys across partitions, and both asc 
and desc) into that
+    // merge with the same kind of LexOrdering the planner builds, and checks 
the output is globally
+    // sorted and complete. It is deterministic coverage of the merge that 
does not depend on an
+    // ordering-reporting Iceberg build (which is why the end-to-end suite's 
merge assertions cancel
+    // on the published Iceberg used in CI).
+    async fn merge_ints(input: Vec<Vec<i32>>, descending: bool) -> Vec<i32> {

Review Comment:
   Picking this thread back up. `merge_ints` and the two `spm_merges_*` tests 
below it still build a `MemorySourceConfig` and merge it, so they never touch 
`IcebergScanExec`, the planner, or the proto, and they still compile and pass 
verbatim on `main`. The mutation I described holds: changing 
`IcebergScanExec::execute` to `let tasks = self.tasks.clone()` so every 
partition reads every file leaves both green.
   
   The three planner tests added alongside them are real coverage of the 
merge-versus-sort decision and of the cap arriving over the proto, so the gap 
is narrower than it was. What is still uncovered is the data path, reading 
actual sorted files through `IcebergScanExec` and checking what comes out.
   
   You asked me for the patch and I dropped it, sorry about that. The shape is 
three real Parquet files written into a tempdir and driven through 
`IcebergScanExec`, ascending, descending, and with duplicate keys across files, 
then the same data through `create_plan` asserting the merge path and the 
`maxFilesPerPartition` sort path agree, plus the null ordering round-tripped 
through the proto in all four direction and null-ordering combinations. It runs 
in about 50ms and the mutation above kills it. Say the word and I will push it 
to your branch.



##########
spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala:
##########
@@ -1022,16 +1023,66 @@ case class CometScanRule(session: SparkSession)
           }
         }
 
+        // If Iceberg reports an ordering, EnsureRequirements may have already 
dropped the Sort
+        // above this scan (it decides that on the vanilla BatchScanExec, 
before Comet converts the
+        // scan). If the native scan cannot guarantee that ordering, reading 
unordered here would
+        // silently return wrong results, so stay on Spark -- its Iceberg 
reader produces the sorted
+        // output it promised. Evaluate the gate exactly once here and stash 
the result on the
+        // metadata; CometIcebergNativeScanExec.outputOrdering and the proto 
serde both read that
+        // stashed value, so the reported order cannot diverge from what 
native advertises.
+        val icebergReportsOrdering: Boolean = 
scanExec.ordering.exists(_.nonEmpty)

Review Comment:
   I went and read the Iceberg side to work out when `orderingHonored` actually 
fires, and I think it costs us the whole native scan in a case where it does 
not have to.
   
   In apache/iceberg#16750, `SparkPartitioningAwareScan.outputOrdering()` 
returns `Spark3Util.toOrdering(table().sortOrder())` whenever 
`isOrderingEnabled()` holds. `SortOrderAnalyzer.canReportOrdering` checks four 
things: the table has a sort order, the sort order has no nested fields, each 
partition key maps to exactly one task group, and every file carries the 
current sort order id. None of them looks at the projection. On the Spark side, 
`V2ScanPartitioningAndOrdering.ordering` resolves those references against 
`relation`, which is the unpruned relation, and unlike the `partitioning` 
branch directly above it there is no `references.subsetOf(d.outputSet)` guard.
   
   So for `SELECT data FROM sorted_table` where the sort key is `id`, 
`scanExec.ordering` comes back non-empty carrying `id`'s exprId, which is not 
in `scanExec.output`. `isIdentityProjected` fails, `reportedOrdering` is `Nil`, 
`orderingHonored` is false, and the whole Iceberg scan goes back to Spark.
   
   I do not think that buys any safety. A required ordering can only reference 
attributes in the scan's output, and `SortOrder.orderingSatisfies` matches 
position by position starting at index 0, so an ordering whose first field is 
an unprojected attribute satisfies nothing above the scan. Spark has dropped no 
`Sort`, and reading unordered would be correct.
   
   Could we keep the scan native when the first reported sort field references 
an attribute outside `scanExec.output`, and report no ordering in that case? 
The general version is to honour the longest prefix of identity-projected 
fields, since that is exactly as far as any satisfiable required ordering can 
reach. Declining is still the right answer when the blocking field is a 
projected UUID, because a required ordering can reach into it, so the UUID path 
would not change.
   
   This matters because it is the common shape. Once `preserve-data-ordering` 
is on, any query on a sorted table that does not select every sort key loses 
the native scan. The follow-up bullet in the PR description still describes 
this case as "we skip the merge and read unordered", which was true before the 
`orderingHonored` gate went in.



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