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]