parthchandra commented on code in PR #5331:
URL: https://github.com/apache/datafusion-comet/pull/5331#discussion_r4160366138
##########
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:
Good catch, fixed. When the sort key isn't one of the selected columns, we
no longer send the whole scan back to Spark. We keep the native scan and simply
report no ordering (nothing above the scan can ask for an order on a column
that isn't selected, so Spark hasn't dropped a sort). We still report the
longest run of sort columns we can honour, and we only fall back to Spark when
a selected column is unsafe — a UUID or a transform — or when we can't read the
table schema. See `CometScanRule.scala:1000-1021` and the new orderingDecision
in `CometIcebergNativeScan.scala:1000`.
##########
native/core/src/execution/planner.rs:
##########
@@ -1844,26 +1845,76 @@ impl PhysicalPlanner {
let metadata_location = common.metadata_location.clone();
let catalog_name = common.catalog_name.clone();
let tasks = parse_file_scan_tasks_from_common(common,
&scan.file_scan_tasks)?;
+ let tasks_len = tasks.len();
let data_file_concurrency_limit =
common.data_file_concurrency_limit as usize;
+ let max_files_per_partition = common.max_files_per_partition
as usize;
- let iceberg_scan = IcebergScanExec::new(
+ // Table sort order Iceberg reported. Empty unless sortMerge
is on and the order
+ // passed the identity gate in CometIcebergNativeScan. The
SortOrder children are
+ // bound references into required_schema, so build the
LexOrdering against it.
+ let ordering: Option<LexOrdering> = if
common.table_sort_orders.is_empty() {
+ None
+ } else {
+ let exprs = common
+ .table_sort_orders
+ .iter()
+ .map(|expr| self.create_sort_expr(expr,
Arc::clone(&required_schema)))
+ .collect::<Result<Vec<PhysicalSortExpr>,
ExecutionError>>()?;
+ LexOrdering::new(exprs)
+ };
+
+ // A per-file-stream k-way merge opens one reader per file at
once. Above the
+ // configured limit we instead read the partition unordered
(bounded by
+ // data_file_concurrency_limit) and sort with a spillable
SortExec, which bounds both
+ // open readers and memory. Both paths still produce sorted
output, so the ordering
+ // Spark eliminated its Sort on is honoured either way. A
limit of 0 (sortMerge
+ // disabled) always takes the sort path.
+ let use_merge = ordering.is_some() && tasks_len <=
max_files_per_partition;
Review Comment:
Fixed. When a float or double column is a sort key but is not the last one,
we no longer do the streaming merge — we read the files unordered and run the
spillable sort instead, which sorts the way Spark does (treating -0.0 and +0.0
as equal). That removes the case where the merge trusted Iceberg's file order
(which keeps -0.0 before +0.0) and produced a wrong result for things like
ORDER BY ... LIMIT. A float/double that is the only key, or the last key, is
still safe, so those still take the merge. See the decision at
`planner.rs:1898-1908`, and three new plan tests: leading float → sort,
trailing float → merge, single float → merge.
##########
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:
Fixed the comment.
##########
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:
Yes please — go ahead and push that test to the branch. I'd rather take
your version that you've already measured than reproduce it. Thank you!
##########
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:
You're right, the old test asserted nothing. Rewritten to check the real
outcome: it now asserts the scan stays native and reports no ordering
(CometIcebergSortMergeReadSuite.scala:643). I also added unit tests for the
decision itself — stays-native for an unselected key, falls back for a selected
UUID/transform, and the "unselected key before a UUID still stays native" case
— starting at line 138.
##########
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:
Fair point. Logged a new issue for this -
https://github.com/apache/datafusion-comet/issues/6524
--
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]