diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index 2bb8fd363d1..e76136c9038 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -471,6 +471,7 @@ jobs: org.apache.comet.CometIcebergEncryptionSuite org.apache.comet.CometIcebergRewriteActionSuite org.apache.comet.CometIcebergWriteActionSuite + org.apache.comet.CometIcebergSortMergeReadSuite org.apache.comet.CometIcebergWriteDetectionSuite org.apache.comet.CometIcebergSystemFunctionSuite org.apache.comet.CometIcebergSystemFunctionExtensionsSuite diff --git a/.github/workflows/pr_build_macos.yml b/.github/workflows/pr_build_macos.yml index 148dde4784e..05bc9bfab06 100644 --- a/.github/workflows/pr_build_macos.yml +++ b/.github/workflows/pr_build_macos.yml @@ -175,6 +175,7 @@ jobs: org.apache.comet.CometIcebergEncryptionSuite org.apache.comet.CometIcebergRewriteActionSuite org.apache.comet.CometIcebergWriteActionSuite + org.apache.comet.CometIcebergSortMergeReadSuite org.apache.comet.CometIcebergWriteDetectionSuite org.apache.comet.CometIcebergSystemFunctionSuite org.apache.comet.CometIcebergSystemFunctionExtensionsSuite diff --git a/docs/source/user-guide/latest/iceberg.md b/docs/source/user-guide/latest/iceberg.md index 4c3abe00216..56aecd37ef5 100644 --- a/docs/source/user-guide/latest/iceberg.md +++ b/docs/source/user-guide/latest/iceberg.md @@ -58,6 +58,26 @@ config `spark.comet.scan.icebergNative.dataFileConcurrencyLimit`. This value def maintain test behavior on Iceberg Java tests without `ORDER BY` clauses, but we suggest increasing it to values between 2 and 8 based on your workload. +When a table is sorted and Iceberg reports its sort order (which requires Iceberg's +`spark.sql.iceberg.planning.preserve-data-ordering`), Comet reads each partition's already-sorted +files as separate streams and merges them, so Spark can drop redundant sorts above the scan. Two +configs tune this: + +- `spark.comet.scan.icebergNative.sortMerge.enabled` (default `true`) turns the streaming merge on + or off. When off, the scan stays native and still honors the reported order, but via a single + unordered read plus a spillable sort rather than the k-way merge. +- `spark.comet.scan.icebergNative.sortMerge.maxFilesPerPartition` (default `64`) caps how many + files in one partition Comet will merge. A merge opens one reader per file at once, so above this + many files Comet reads the partition unordered and sorts it with a spillable sort instead, which + bounds concurrently-open readers. Both paths produce sorted output. + +Sort elimination only applies when the sort key columns are in the query's projection. If a query +does not select a leading sort key (for example `SELECT data FROM sorted_table` where the table is +sorted by `id`), Comet still reads the table with the native scan but reports no ordering, so Spark +keeps any sort it would otherwise have. This is not a fallback to Spark's reader -- only sort +elimination is given up. It is safe because a required ordering can only reference projected columns, +so Spark never eliminated a sort that depended on the unprojected key. + ### Supported features The native Iceberg reader supports the following features: @@ -196,6 +216,12 @@ The following scenarios will fall back to the JVM Iceberg reader: or renamed. The native reader cannot yet match such a field to data files written before the change. The check uses the table's schema history, so the fallback stays after those files are rewritten +- 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 + while Spark maps UUID to a string, so the two orders can disagree; Comet declines the ordering + and falls back to Spark Writes are not covered by this list. By default Iceberg writes use Spark's own writer; see [Iceberg Writes](iceberg-writes.md) for the experimental native writer and when it applies. diff --git a/docs/source/user-guide/latest/tuning/scans.md b/docs/source/user-guide/latest/tuning/scans.md index 40bc82634ec..6a1c946736e 100644 --- a/docs/source/user-guide/latest/tuning/scans.md +++ b/docs/source/user-guide/latest/tuning/scans.md @@ -67,6 +67,13 @@ worked example and further discussion. ## Iceberg Scan Tuning Comet's native Iceberg scan (`spark.comet.scan.icebergNative.enabled`, enabled by default) reads each -task's data files one at a time by default. For tables with many small files or high-latency storage, -increase `spark.comet.scan.icebergNative.dataFileConcurrencyLimit` (default `1`; values of 2–8 are -suggested) to overlap I/O across files at the cost of extra memory. +task's data files one at a time by default on the unordered read path. For tables with many small files +or high-latency storage, increase `spark.comet.scan.icebergNative.dataFileConcurrencyLimit` (default +`1`; values of 2–8 are suggested) to overlap I/O across files at the cost of extra memory. + +This limit only bounds the unordered read. When reading a sorted Iceberg table that reports its +ordering (`spark.sql.iceberg.planning.preserve-data-ordering`), each file becomes its own partition +under the streaming merge, so a task can hold up to `spark.comet.scan.icebergNative.sortMerge.maxFilesPerPartition` +(default `64`) readers open at once, regardless of `dataFileConcurrencyLimit`. To cap scan memory on a +sorted table, lower `sortMerge.maxFilesPerPartition` (above it the scan reads unordered and sorts with a +spillable sort instead of merging). diff --git a/native/core/src/execution/operators/iceberg_scan.rs b/native/core/src/execution/operators/iceberg_scan.rs index 9f465329ca0..625c43d6dfb 100644 --- a/native/core/src/execution/operators/iceberg_scan.rs +++ b/native/core/src/execution/operators/iceberg_scan.rs @@ -20,7 +20,7 @@ use std::collections::HashMap; use std::fmt; use std::pin::Pin; -use std::sync::Arc; +use std::sync::{Arc, Mutex, OnceLock}; use std::task::{Context, Poll}; use arrow::array::{ArrayRef, RecordBatch, RecordBatchOptions}; @@ -29,7 +29,7 @@ use datafusion::common::tree_node::TreeNodeRecursion; use datafusion::common::{DataFusionError, Result as DFResult}; use datafusion::execution::{RecordBatchStream, SendableRecordBatchStream, TaskContext}; use datafusion::physical_expr::expressions::Column; -use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr}; +use datafusion::physical_expr::{EquivalenceProperties, LexOrdering, PhysicalExpr}; use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType}; use datafusion::physical_plan::metrics::{ BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricsSet, @@ -38,7 +38,7 @@ use datafusion::physical_plan::{ DisplayAs, DisplayFormatType, ExecutionPlan, Partitioning, PlanProperties, }; use futures::{Stream, StreamExt, TryStreamExt}; -use iceberg::arrow::ScanMetrics; +use iceberg::arrow::{ArrowReader, ScanMetrics}; use iceberg::io::FileIO; use iceberg::Runtime as IcebergRuntime; use iceberg::{Error, ErrorKind}; @@ -83,6 +83,36 @@ pub struct IcebergScanExec { tasks: Vec, /// Number of data files to read concurrently data_file_concurrency_limit: usize, + /// FileIO (and, for S3, the JVM credential bridge behind it) built once at plan time and shared + /// across partitions. FileIO is cheap to clone (Arc-backed), so each `execute` clones this + /// rather than rebuilding the storage factory + credential bridge. This matters in the ordered + /// path, where the scan is one partition per file and `execute` is called once per file. + file_io: FileIO, + /// Table sort order Iceberg reported, translated against `output_schema`. `Some` makes this a + /// multi-partition scan: one sorted stream per task, which a SortPreservingMergeExec above + /// merges back into one sorted partition. It is also advertised in `plan_properties`. `None` + /// keeps the old single-partition unordered read (all tasks streamed together). + /// + /// Concurrency note: in the ordered path each partition reads exactly one task, so + /// `data_file_concurrency_limit` no longer bounds cross-file concurrency; instead the wrapping + /// SortPreservingMergeExec drives one reader per file to merge them. The planner only takes + /// this ordered/merge path when a partition has at most `sortMerge.maxFilesPerPartition` files; + /// above that it reads the partition unordered and wraps a spillable `SortExec`, so this fan-out + /// (one reader per file) is capped rather than unbounded. A finer, bound-driven admission scheme + /// that reads a subset of files at a time using per-file min/max is tracked in #5343. + /// `data_file_concurrency_limit` still bounds delete-file stats and the unordered path. + ordering: Option, + /// One ArrowReader shared across this scan's partitions, built lazily on the first `execute` + /// (it needs the session batch size). In the ordered path each partition reads one file, so + /// sharing the reader shares its `CachingDeleteFileLoader`: a delete file that applies to the + /// whole partition is downloaded and parsed once, not once per data file (#6524). `read()` + /// creates fresh `ScanMetrics` per call, so per-partition metrics stay separate, and the reader + /// is Arc-backed and cheap to clone. + reader: OnceLock, + /// Delete-file sizes (path -> bytes) stat'd once and shared across partitions, so the HEAD + /// request for a delete file shared by the partition is issued once rather than once per data + /// file (#6524). Guarded by a std Mutex; the lock is never held across an `.await`. + delete_file_sizes: Arc>>, /// Metrics metrics: ExecutionPlanMetricsSet, } @@ -95,12 +125,36 @@ impl IcebergScanExec { catalog_name: String, tasks: Vec, data_file_concurrency_limit: usize, + ordering: Option, ) -> Result { let output_schema = schema; - let plan_properties = Self::compute_properties(Arc::clone(&output_schema), 1); + // With an ordering, read each task as its own sorted stream on its own partition, so the + // SortPreservingMergeExec above can merge them. Without one, keep the single partition that + // reads every task. Comet only drives execute(0), so a multi-partition leaf with no merge + // above it would read only the first task. + let num_partitions = if ordering.is_some() { + tasks.len().max(1) + } else { + 1 + }; + let plan_properties = Self::compute_properties( + Arc::clone(&output_schema), + num_partitions, + ordering.as_ref(), + ); let metrics = ExecutionPlanMetricsSet::new(); + // Build FileIO (and the S3 credential bridge) once here rather than per `execute`. In the + // ordered path `execute` is called once per file, so rebuilding it there would repeat the + // JNI/reflection credential-bridge construction for every file in the partition. + let file_io = load_file_io( + &catalog_properties, + &metadata_location, + &catalog_name, + AccessMode::Read, + )?; + Ok(Self { metadata_location, output_schema, @@ -109,13 +163,28 @@ impl IcebergScanExec { catalog_name, tasks, data_file_concurrency_limit, + file_io, + ordering, + reader: OnceLock::new(), + delete_file_sizes: Arc::new(Mutex::new(HashMap::new())), metrics, }) } - fn compute_properties(schema: SchemaRef, num_partitions: usize) -> Arc { + fn compute_properties( + schema: SchemaRef, + num_partitions: usize, + ordering: Option<&LexOrdering>, + ) -> Arc { + let eq_properties = match ordering { + Some(lex) => EquivalenceProperties::new_with_orderings( + Arc::clone(&schema), + std::iter::once(lex.iter().cloned()), + ), + None => EquivalenceProperties::new(Arc::clone(&schema)), + }; Arc::new(PlanProperties::new( - EquivalenceProperties::new(schema), + eq_properties, Partitioning::UnknownPartitioning(num_partitions), EmissionType::Incremental, Boundedness::Bounded, @@ -158,10 +227,18 @@ impl ExecutionPlan for IcebergScanExec { fn execute( &self, - _partition: usize, + partition: usize, context: Arc, ) -> DFResult { - self.execute_with_tasks(self.tasks.clone(), context) + // With a reported ordering this is a multi-partition operator: partition `i` reads only + // task `i` as its own sorted stream, and the SortPreservingMergeExec above merges them. + // Without one, the single partition reads every task together (legacy unordered path). + let tasks = if self.ordering.is_some() { + self.tasks.get(partition).cloned().into_iter().collect() + } else { + self.tasks.clone() + }; + self.execute_with_tasks(tasks, partition, context) } fn metrics(&self) -> Option { @@ -175,51 +252,63 @@ impl IcebergScanExec { fn execute_with_tasks( &self, tasks: Vec, + partition: usize, context: Arc, ) -> DFResult { 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(); let batch_size = context.session_config().batch_size(); - let metrics = IcebergScanMetrics::new(&self.metrics); + let metrics = IcebergScanMetrics::new(&self.metrics, partition); metrics.num_splits.add(tasks.len()); // Fill delete-file sizes as the first step of the task stream so the stats run on the // iceberg runtime alongside the reads, not on the calling executor thread (see - // fill_delete_file_sizes). + // fill_delete_file_sizes). The size cache is shared across partitions so a delete file the + // whole partition shares is stat'd once, not once per data file (#6524). let fill_io = file_io.clone(); let concurrency_limit = self.data_file_concurrency_limit; + let delete_sizes = Arc::clone(&self.delete_file_sizes); let task_stream = futures::stream::once(async move { let mut tasks = tasks; - Self::fill_delete_file_sizes(&mut tasks, &fill_io, concurrency_limit).await?; + Self::fill_delete_file_sizes(&mut tasks, &fill_io, concurrency_limit, &delete_sizes) + .await?; Ok::<_, Error>(futures::stream::iter(tasks.into_iter().map(Ok::<_, Error>))) }) .try_flatten() .boxed(); - // iceberg-rust's ArrowReader spawns IO/CPU work onto an iceberg::Runtime, which only needs - // a tokio handle. execute() runs on the JVM-called thread outside any tokio context, so we - // enter Comet's global runtime to capture its handle (this is where the stream is later - // polled). Capturing the handle rather than borrowing the runtime keeps it tear-downable - // via release_runtime. - let iceberg_runtime = { - let handle = get_runtime(); - let _guard = handle.enter(); - IcebergRuntime::try_current().map_err(|e| { - DataFusionError::Execution(format!("Failed to build Iceberg runtime: {e}")) - })? + // Build the ArrowReader once and share it across the scan's partitions, so the per-file + // streams share one CachingDeleteFileLoader (see the `reader` field, #6524). The first + // caller builds it; a later caller reuses that one, so the delete cache is shared. read() + // makes fresh ScanMetrics per call, so per-partition metrics stay separate. + let reader = match self.reader.get() { + Some(reader) => reader.clone(), + None => { + // iceberg-rust's ArrowReader spawns IO/CPU work onto an iceberg::Runtime, which only + // needs a tokio handle. execute() runs on the JVM-called thread outside any tokio + // context, so enter Comet's global runtime to capture its handle (this is where the + // stream is later polled). Capturing the handle rather than borrowing the runtime + // keeps it tear-downable via release_runtime. + let iceberg_runtime = { + let handle = get_runtime(); + let _guard = handle.enter(); + IcebergRuntime::try_current().map_err(|e| { + DataFusionError::Execution(format!("Failed to build Iceberg runtime: {e}")) + })? + }; + let built = iceberg::arrow::ArrowReaderBuilder::new(file_io, iceberg_runtime) + .with_batch_size(batch_size) + .with_data_file_concurrency_limit(self.data_file_concurrency_limit) + .with_row_selection_enabled(true) + .with_metadata_size_hint(512 * 1024) // Same as DataFusion's default + .build(); + // First caller wins; `set` is a no-op (returns Err) if another partition already set + // it, so every partition converges on one shared reader. + let _ = self.reader.set(built); + self.reader.get().expect("reader was just set").clone() + } }; - let reader = iceberg::arrow::ArrowReaderBuilder::new(file_io, iceberg_runtime) - .with_batch_size(batch_size) - .with_data_file_concurrency_limit(self.data_file_concurrency_limit) - .with_row_selection_enabled(true) - .with_metadata_size_hint(512 * 1024) // Same as DataFusion's default - .build(); // Pass all tasks to iceberg-rust at once to utilize its flatten_unordered // parallelization, avoiding overhead of single-task streams @@ -259,8 +348,9 @@ impl IcebergScanExec { tasks: &mut [FileScanTask], file_io: &FileIO, concurrency_limit: usize, + size_cache: &Mutex>, ) -> Result<(), Error> { - use datafusion::common::{HashMap, HashSet}; + use datafusion::common::HashSet; use futures::TryStreamExt; // Dedup: the JVM pools delete-file lists, not individual files, so the same delete file @@ -286,13 +376,26 @@ impl IcebergScanExec { return Ok(()); } + // Only stat delete files not already sized on an earlier partition of this scan. The size + // cache is shared across partitions (#6524), so a delete file the whole partition shares is + // stat'd once, not once per data file. Snapshot the missing paths under the lock, then + // release it before the async stats below. + let to_stat: Vec = { + let cache = size_cache.lock().unwrap(); + needed + .iter() + .filter(|path| !cache.contains_key(path.as_str())) + .cloned() + .collect() + }; + // Bound the in-flight stats to match the downstream read concurrency (iceberg-rust uses // try_buffer_unordered at the same limit for delete-file loads). An unbounded fan-out // would burst N HEAD requests at once for no gain, since the reads are throttled anyway. // Guaranteed > 0 by COMET_ICEBERG_DATA_FILE_CONCURRENCY_LIMIT; buffer_unordered(0) would // never poll. debug_assert!(concurrency_limit > 0); - let sizes: Vec<(String, u64)> = futures::stream::iter(needed.into_iter().map(|path| { + let sizes: Vec<(String, u64)> = futures::stream::iter(to_stat.into_iter().map(|path| { let file_io = file_io.clone(); async move { let size = file_io @@ -325,7 +428,18 @@ impl IcebergScanExec { .try_collect() .await?; - let size_map: HashMap = sizes.into_iter().collect(); + // Record newly-stat'd sizes in the shared cache, then build the size map for every delete + // file this partition needs (cached + just stat'd). The lock is not held across any await. + let size_map: std::collections::HashMap = { + let mut cache = size_cache.lock().unwrap(); + for (path, size) in sizes { + cache.insert(path, size); + } + needed + .iter() + .filter_map(|path| cache.get(path.as_str()).map(|&size| (path.clone(), size))) + .collect() + }; // iceberg-rust 665c64e made `FileScanTask::deletes` a private field exposed only through // a read-only accessor, so a task's delete files can no longer be sized in place. Rebuild // each task that carries deletes with sized copies; tasks without deletes are untouched. @@ -392,11 +506,11 @@ struct IcebergScanMetrics { } impl IcebergScanMetrics { - fn new(metrics: &ExecutionPlanMetricsSet) -> Self { + fn new(metrics: &ExecutionPlanMetricsSet, partition: usize) -> Self { Self { - baseline: BaselineMetrics::new(metrics, 0), - num_splits: MetricBuilder::new(metrics).counter("num_splits", 0), - bytes_scanned: MetricBuilder::new(metrics).counter("bytes_scanned", 0), + baseline: BaselineMetrics::new(metrics, partition), + num_splits: MetricBuilder::new(metrics).counter("num_splits", partition), + bytes_scanned: MetricBuilder::new(metrics).counter("bytes_scanned", partition), } } } @@ -614,6 +728,7 @@ mod tests { use iceberg_storage_opendal::OpenDalStorageFactory; use super::IcebergScanExec; + use datafusion::physical_plan::ExecutionPlan; fn fs_file_io() -> FileIO { FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Fs)).build() @@ -714,7 +829,13 @@ mod tests { missing.to_str().unwrap(), )])]; - let result = IcebergScanExec::fill_delete_file_sizes(&mut tasks, &fs_file_io(), 4).await; + let result = IcebergScanExec::fill_delete_file_sizes( + &mut tasks, + &fs_file_io(), + 4, + &std::sync::Mutex::new(std::collections::HashMap::new()), + ) + .await; assert!( result.is_err(), @@ -734,9 +855,14 @@ mod tests { f.flush().unwrap(); let mut tasks = vec![task_with_deletes(vec![delete_file(path.to_str().unwrap())])]; - IcebergScanExec::fill_delete_file_sizes(&mut tasks, &fs_file_io(), 4) - .await - .unwrap(); + IcebergScanExec::fill_delete_file_sizes( + &mut tasks, + &fs_file_io(), + 4, + &std::sync::Mutex::new(std::collections::HashMap::new()), + ) + .await + .unwrap(); assert_eq!(tasks[0].deletes()[0].file_size_in_bytes, bytes.len() as u64); } @@ -751,7 +877,13 @@ mod tests { std::fs::File::create(&path).unwrap(); // 0 bytes on disk let mut tasks = vec![task_with_deletes(vec![delete_file(path.to_str().unwrap())])]; - let result = IcebergScanExec::fill_delete_file_sizes(&mut tasks, &fs_file_io(), 4).await; + let result = IcebergScanExec::fill_delete_file_sizes( + &mut tasks, + &fs_file_io(), + 4, + &std::sync::Mutex::new(std::collections::HashMap::new()), + ) + .await; assert!( result.is_err(), @@ -764,9 +896,149 @@ mod tests { #[tokio::test] async fn fill_delete_file_sizes_noop_without_deletes() { let mut tasks = vec![task_with_deletes(vec![])]; - IcebergScanExec::fill_delete_file_sizes(&mut tasks, &fs_file_io(), 4) - .await - .unwrap(); + IcebergScanExec::fill_delete_file_sizes( + &mut tasks, + &fs_file_io(), + 4, + &std::sync::Mutex::new(std::collections::HashMap::new()), + ) + .await + .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 { + 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, + ) -> 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>, descending: bool) -> Vec { + use arrow::array::Int32Array; + use arrow::array::RecordBatch; + use arrow::compute::SortOptions; + use arrow::datatypes::{DataType, Field, Schema}; + use datafusion::datasource::memory::MemorySourceConfig; + use datafusion::physical_expr::expressions::Column; + use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr}; + use datafusion::physical_plan::sorts::sort_preserving_merge::SortPreservingMergeExec; + use datafusion::prelude::SessionContext; + use futures::StreamExt; + + let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)])); + let partitions: Vec> = input + .into_iter() + .map(|vals| { + vec![RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from(vals))], + ) + .unwrap()] + }) + .collect(); + let source = + MemorySourceConfig::try_new_exec(&partitions, Arc::clone(&schema), None).unwrap(); + + let ordering = LexOrdering::new(vec![PhysicalSortExpr { + expr: Arc::new(Column::new("id", 0)), + options: SortOptions { + descending, + nulls_first: false, + }, + }]) + .unwrap(); + let spm = SortPreservingMergeExec::new(ordering, source); + assert_eq!(spm.properties().partitioning.partition_count(), 1); + + let ctx = SessionContext::new(); + let mut stream = spm.execute(0, ctx.task_ctx()).unwrap(); + let mut got = Vec::new(); + while let Some(batch) = stream.next().await { + let batch = batch.unwrap(); + let col = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + (0..col.len()).for_each(|i| got.push(col.value(i))); + } + got + } + + #[tokio::test] + async fn spm_merges_ascending_sorted_partitions() { + // Two individually-sorted partitions with duplicate keys across them. + let got = merge_ints(vec![vec![1, 3, 3, 5, 8], vec![2, 3, 6, 7]], false).await; + let mut expected = vec![1, 3, 3, 5, 8, 2, 3, 6, 7]; + expected.sort_unstable(); + assert_eq!( + got, expected, + "ascending merge must be globally sorted and complete" + ); + } + + #[tokio::test] + async fn spm_merges_descending_sorted_partitions() { + let got = merge_ints(vec![vec![8, 5, 3, 3, 1], vec![7, 6, 3, 2]], true).await; + let mut expected = vec![8, 5, 3, 3, 1, 7, 6, 3, 2]; + expected.sort_unstable_by(|a, b| b.cmp(a)); + assert_eq!( + got, expected, + "descending merge must be globally sorted and complete" + ); } fn from_hex(s: &str) -> Vec { diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index b88b477931d..85fee63d1fd 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -79,6 +79,7 @@ use datafusion::{ limit::LocalLimitExec, projection::ProjectionExec, sorts::sort::SortExec, + sorts::sort_preserving_merge::SortPreservingMergeExec, ChildrenPropertiesMode, ExecutionPlan, ExecutionPlanProperties, ReplaceChildrenOptions, }, prelude::SessionContext, @@ -2159,26 +2160,102 @@ 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; + + // Table sort order Iceberg reported. Empty unless the order passed the identity gate + // in CometScanRule; the serde writes it regardless of sortMerge.enabled. When + // sortMerge is disabled the ordering is still reported (enabled=false only sets + // max_files_per_partition to 0, so the branch below takes the SortExec path) -- an + // empty list there would mean an unordered read under a Sort that Spark already + // dropped. The SortOrder children are bound references into required_schema, so build + // the LexOrdering against it. + let ordering: Option = 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::, 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. + // + // A float/double sort key that PRECEDES another key also forces the sort path: the + // merge trusts each file to already be sorted under Spark's comparator, but Iceberg + // orders files with its own comparator, which distinguishes -0.0 from +0.0 while + // Spark's create_sort_expr treats them equal. So a file Iceberg considers sorted -- + // e.g. (-0.0, 2) then (+0.0, 1) -- is out of order under the composite Spark + // comparator, and merging would preserve that wrong order (which can change results, + // e.g. a TakeOrderedAndProject that trusts the advertised ordering). The spillable + // SortExec re-sorts under Spark's comparator and is correct. A trailing float key is + // safe: signed-zero rows compare equal in Spark, so any order among them already + // satisfies the required ordering. A field whose type cannot be resolved is treated + // as unsafe (force the sort), so we never merge on an unverified key. + let ordering_merge_safe = ordering.as_ref().is_some_and(|lex| { + let last = lex.len().saturating_sub(1); + lex.iter().enumerate().all(|(i, sort_expr)| { + i == last + || matches!( + sort_expr.expr.data_type(&required_schema), + Ok(dt) if !matches!(dt, DataType::Float32 | DataType::Float64) + ) + }) + }); + let use_merge = ordering_merge_safe && tasks_len <= max_files_per_partition; + + // Only the merge path is multi-partition (one sorted stream per file). The sort + // fallback and the no-ordering path read a single unordered partition. + let scan_ordering = if use_merge { ordering.clone() } else { None }; - let iceberg_scan = IcebergScanExec::new( + let iceberg_scan: Arc = Arc::new(IcebergScanExec::new( metadata_location, required_schema, catalog_properties, catalog_name, tasks, data_file_concurrency_limit, - )?; + scan_ordering, + )?); - Ok(( - vec![], - vec![], - Arc::new(SparkPlan::new( - spark_plan.plan_id, - Arc::new(iceberg_scan), - vec![], - )), - )) + // Metrics of the underlying scan roll up via additional_native_plans so num_splits / + // bytes_scanned still surface under the single Spark scan node. + let result_plan = match ordering { + Some(lex) if use_merge => { + let spm: Arc = Arc::new( + SortPreservingMergeExec::new(lex, Arc::clone(&iceberg_scan)) + .with_round_robin_repartition(false), + ); + SparkPlan::new_with_additional( + spark_plan.plan_id, + spm, + vec![], + vec![iceberg_scan], + ) + } + Some(lex) => { + // Too many files to merge: single unordered read + spillable SortExec. + let sort: Arc = + Arc::new(SortExec::new(lex, Arc::clone(&iceberg_scan))); + SparkPlan::new_with_additional( + spark_plan.plan_id, + sort, + vec![], + vec![iceberg_scan], + ) + } + None => SparkPlan::new(spark_plan.plan_id, iceberg_scan, vec![]), + }; + + Ok((vec![], vec![], Arc::new(result_plan))) } OpStruct::ContribScan(contrib) => { // Extension point for optional, out-of-tree contrib scans (Delta, Lance, ...). The @@ -8027,6 +8104,255 @@ mod tests { ); } + /// Builds an IcebergScan Operator with `num_files` file-scan tasks and a reported identity + /// ordering on a single Int column, so create_plan takes the sort-merge path. + fn iceberg_scan_op_with_ordering(num_files: usize, max_files_per_partition: u32) -> Operator { + let schema_json = serde_json::to_string( + &iceberg::spec::Schema::builder() + .with_schema_id(0) + .with_fields(vec![iceberg::spec::NestedField::required( + 1, + "id", + iceberg::spec::Type::Primitive(iceberg::spec::PrimitiveType::Int), + ) + .into()]) + .build() + .expect("schema"), + ) + .expect("serialize schema"); + + let int_type = spark_expression::DataType { + type_id: 3, // Int32 + type_info: None, + }; + + let required_schema = vec![spark_operator::SparkStructField { + name: "id".to_string(), + data_type: Some(int_type.clone()), + nullable: false, + metadata: Default::default(), + }]; + + // id ASC NULLS FIRST, as Expr { SortOrder { child: Bound(0) } }. + let sort_child = Expr { + expr_struct: Some(Bound(spark_expression::BoundReference { + index: 0, + datatype: Some(int_type), + })), + query_context: None, + expr_id: None, + }; + let sort_order = Expr { + expr_struct: Some(ExprStruct::SortOrder(Box::new( + spark_expression::SortOrder { + child: Some(Box::new(sort_child)), + direction: 0, // Ascending + null_ordering: 0, // NullsFirst + }, + ))), + query_context: None, + expr_id: None, + }; + + let common = spark_operator::IcebergScanCommon { + schema_pool: vec![schema_json], + project_field_ids_pool: vec![spark_operator::ProjectFieldIdList { field_ids: vec![1] }], + required_schema, + table_sort_orders: vec![sort_order], + max_files_per_partition, + ..Default::default() + }; + + let task = spark_operator::IcebergFileScanTask { + data_file_path: "file:///tmp/data.parquet".to_string(), + file_size_in_bytes: 100, + schema_idx: 0, + project_field_ids_idx: 0, + ..Default::default() + }; + + Operator { + plan_id: 1, + sql_text_pool: vec![], + children: vec![], + op_struct: Some(OpStruct::IcebergScan(spark_operator::IcebergScan { + common: Some(common), + file_scan_tasks: vec![task; num_files], + })), + } + } + + // At or below the max-files-per-partition limit, a reported ordering makes the scan + // multi-partition and the planner wraps it in a SortPreservingMergeExec. + #[test] + fn iceberg_scan_merges_when_files_within_limit() { + let op = iceberg_scan_op_with_ordering(4, 64); + let planner = PhysicalPlanner::default(); + let (_, _, plan) = planner.create_plan(&op, &mut vec![], 1).unwrap(); + assert_eq!("SortPreservingMergeExec", plan.native_plan.name()); + } + + // Above the limit, merging would open too many readers at once, so the planner falls back to a + // single unordered read wrapped in a spillable SortExec (still sorted output). + #[test] + fn iceberg_scan_falls_back_to_sort_above_file_limit() { + let op = iceberg_scan_op_with_ordering(65, 64); + let planner = PhysicalPlanner::default(); + let (_, _, plan) = planner.create_plan(&op, &mut vec![], 1).unwrap(); + assert_eq!("SortExec", plan.native_plan.name()); + } + + // The limit that decides merge-vs-sort comes from the proto (common.max_files_per_partition), + // not a compiled-in constant, so a low limit must flip the same 3-file scan from merge to sort. + // This proves the value is actually read from the plan: with max_files_per_partition = 2, a + // 3-file partition exceeds the limit and takes the SortExec path. + #[test] + fn iceberg_scan_honors_max_files_per_partition_from_proto() { + let merge = iceberg_scan_op_with_ordering(3, 3); + let planner = PhysicalPlanner::default(); + let (_, _, plan) = planner.create_plan(&merge, &mut vec![], 1).unwrap(); + assert_eq!("SortPreservingMergeExec", plan.native_plan.name()); + + let sort = iceberg_scan_op_with_ordering(3, 2); + let (_, _, plan) = planner.create_plan(&sort, &mut vec![], 1).unwrap(); + assert_eq!("SortExec", plan.native_plan.name()); + } + + // Builds an IcebergScan Operator whose sort order is `col_type_ids` (proto DataTypeId: 3=Int32, + // 6=Float64/Double), all ASC, over columns c0, c1, .... Lets the merge-vs-sort decision be + // exercised for a float key at any position. + fn iceberg_scan_op_with_typed_ordering( + col_type_ids: &[i32], + num_files: usize, + max_files_per_partition: u32, + ) -> Operator { + let iceberg_fields: Vec<_> = col_type_ids + .iter() + .enumerate() + .map(|(i, &tid)| { + let prim = if tid == 6 { + iceberg::spec::PrimitiveType::Double + } else { + iceberg::spec::PrimitiveType::Int + }; + iceberg::spec::NestedField::required( + (i + 1) as i32, + format!("c{i}"), + iceberg::spec::Type::Primitive(prim), + ) + .into() + }) + .collect(); + let schema_json = serde_json::to_string( + &iceberg::spec::Schema::builder() + .with_schema_id(0) + .with_fields(iceberg_fields) + .build() + .expect("schema"), + ) + .expect("serialize schema"); + + let required_schema: Vec<_> = col_type_ids + .iter() + .enumerate() + .map(|(i, &tid)| spark_operator::SparkStructField { + name: format!("c{i}"), + data_type: Some(spark_expression::DataType { + type_id: tid, + type_info: None, + }), + nullable: false, + metadata: Default::default(), + }) + .collect(); + + // c0 ASC, c1 ASC, ... each a Bound reference at its own index, NULLS FIRST. + let sort_orders: Vec = col_type_ids + .iter() + .enumerate() + .map(|(i, &tid)| Expr { + expr_struct: Some(ExprStruct::SortOrder(Box::new( + spark_expression::SortOrder { + child: Some(Box::new(Expr { + expr_struct: Some(Bound(spark_expression::BoundReference { + index: i as i32, + datatype: Some(spark_expression::DataType { + type_id: tid, + type_info: None, + }), + })), + query_context: None, + expr_id: None, + })), + direction: 0, // Ascending + null_ordering: 0, // NullsFirst + }, + ))), + query_context: None, + expr_id: None, + }) + .collect(); + + let field_ids: Vec = (1..=col_type_ids.len() as i32).collect(); + let common = spark_operator::IcebergScanCommon { + schema_pool: vec![schema_json], + project_field_ids_pool: vec![spark_operator::ProjectFieldIdList { field_ids }], + required_schema, + table_sort_orders: sort_orders, + max_files_per_partition, + ..Default::default() + }; + + let task = spark_operator::IcebergFileScanTask { + data_file_path: "file:///tmp/data.parquet".to_string(), + file_size_in_bytes: 100, + schema_idx: 0, + project_field_ids_idx: 0, + ..Default::default() + }; + + Operator { + plan_id: 1, + sql_text_pool: vec![], + children: vec![], + op_struct: Some(OpStruct::IcebergScan(spark_operator::IcebergScan { + common: Some(common), + file_scan_tasks: vec![task; num_files], + })), + } + } + + // A float/double sort key that PRECEDES another key is not merge-safe (Iceberg orders -0.0 < + // +0.0 while Spark treats them equal), so the planner routes it to the spillable SortExec even + // under the cap. Regression for the composite floating-point ordering P1 (PR #5331). + #[test] + fn iceberg_scan_uses_sort_for_leading_float_composite_ordering() { + let op = iceberg_scan_op_with_typed_ordering(&[6, 3], 3, 64); // (DOUBLE, INT) + let planner = PhysicalPlanner::default(); + let (_, _, plan) = planner.create_plan(&op, &mut vec![], 1).unwrap(); + assert_eq!("SortExec", plan.native_plan.name()); + } + + // A trailing float key is merge-safe: signed-zero rows are Spark-equal, so any order among them + // already satisfies the required ordering. It must still take the k-way merge under the cap. + #[test] + fn iceberg_scan_merges_for_trailing_float_ordering() { + let op = iceberg_scan_op_with_typed_ordering(&[3, 6], 3, 64); // (INT, DOUBLE) + let planner = PhysicalPlanner::default(); + let (_, _, plan) = planner.create_plan(&op, &mut vec![], 1).unwrap(); + assert_eq!("SortPreservingMergeExec", plan.native_plan.name()); + } + + // A single float key is the last (and only) key, so it is merge-safe and must merge under the + // cap -- guards the `i == last` carve-out from being dropped. + #[test] + fn iceberg_scan_merges_for_single_float_ordering() { + let op = iceberg_scan_op_with_typed_ordering(&[6], 3, 64); // (DOUBLE) + let planner = PhysicalPlanner::default(); + let (_, _, plan) = planner.create_plan(&op, &mut vec![], 1).unwrap(); + assert_eq!("SortPreservingMergeExec", plan.native_plan.name()); + } + /// The planner promises DataFusion a native UDF's return type with every nested field nullable, /// whatever the kernel reports, and every batch the UDF produces has that type. `make_struct_c` /// reports its struct's field `a` as non-nullable for a registration that declares it nullable. diff --git a/native/proto/src/proto/operator.proto b/native/proto/src/proto/operator.proto index 97c65f0da36..9b632ee408e 100644 --- a/native/proto/src/proto/operator.proto +++ b/native/proto/src/proto/operator.proto @@ -353,11 +353,23 @@ message IcebergScanCommon { // cannot fold those entries together because each carries a distinct content offset. repeated string delete_file_path_pool = 16; + // Max files in one Spark partition the native scan will k-way merge to preserve sort order. A + // merge opens one reader per file at once, so above this the scan reads the partition unordered + // and sorts with a spillable SortExec instead (still sorted output). 0 means never merge -- used + // when sortMerge.enabled is off, where the order is still reported and honoured via the sort. + uint32 max_files_per_partition = 18; + // Iceberg's SparkScan.hashCode() (folds in pushed filters, snapshot, branch, and read schema). // A self-join/self-merge can read the same table (same metadata_location) twice with different // projections or filters, so PlanDataInjector keys planning data by metadata_location, // scan_hash_code and the operator's plan_id rather than metadata_location alone. int32 scan_hash_code = 15; + + // Table sort order Iceberg reported (SupportsReportOrdering), against required_schema. Each Expr + // wraps a SortOrder (child, direction, null_ordering). When set, the native scan reads each file + // as its own sorted stream and wraps the scan in a SortPreservingMergeExec, so each Spark + // partition comes out sorted. Empty means no reported order: read unordered, as before. + repeated spark.spark_expression.Expr table_sort_orders = 17; } message IcebergScan { diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index a187f8fa5f8..97206217db4 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -121,6 +121,36 @@ object CometConf extends ShimCometConf { .booleanConf .createWithDefault(true) + val COMET_ICEBERG_SORT_MERGE_ENABLED: ConfigEntry[Boolean] = + conf("spark.comet.scan.icebergNative.sortMerge.enabled") + .category(CATEGORY_SCAN) + .doc("Whether the native Iceberg scan performs a per-partition streaming k-way merge of " + + "already-sorted files. When enabled and Iceberg reports an ordering (requires Iceberg's " + + "spark.sql.iceberg.planning.preserve-data-ordering), each Spark partition reads its " + + "files as separate sorted streams merged into one sorted output. When disabled, the scan " + + "stays native and the ordering is still reported and honoured, but via a single " + + "unordered read wrapped in a spillable sort instead of the k-way merge (equivalent to " + + "setting maxFilesPerPartition to 0). Either way the ordering is surfaced to Spark so " + + "redundant sorts are eliminated.") + .booleanConf + .createWithDefault(true) + + val COMET_ICEBERG_SORT_MERGE_MAX_FILES_PER_PARTITION: ConfigEntry[Int] = + conf("spark.comet.scan.icebergNative.sortMerge.maxFilesPerPartition") + .category(CATEGORY_SCAN) + .doc( + "The maximum number of files in a single partition that the native Iceberg scan will " + + "k-way merge to preserve sort order. A merge opens one reader per file at once, so " + + "above this many files the scan instead reads the partition unordered and sorts it " + + "with a spillable sort, bounding concurrently-open readers. Both paths produce sorted " + + "output; this only trades merge for a full sort on partitions with many files.") + .intConf + .checkValue( + v => v >= 0, + "Max files per partition must be non-negative (0 means never merge; the reported " + + "ordering is honoured with a spillable sort instead)") + .createWithDefault(64) + val COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.write.iceberg.splitOperator.enabled") .category(CATEGORY_TESTING) @@ -145,9 +175,13 @@ object CometConf extends ShimCometConf { conf("spark.comet.scan.icebergNative.dataFileConcurrencyLimit") .category(CATEGORY_SCAN) .doc( - "The number of Iceberg data files to read concurrently within a single task. " + - "Higher values improve throughput for tables with many small files by overlapping " + - "I/O latency, but increase memory usage. Values between 2 and 8 are suggested.") + "The number of Iceberg data files to read concurrently within a single task on the " + + "unordered read path. Higher values improve throughput for tables with many small " + + "files by overlapping I/O latency, but increase memory usage. Values between 2 and 8 " + + "are suggested. This does NOT bound the sorted (merge) path: when Iceberg reports an " + + "ordering, each file is its own partition under the merge, so up to " + + "sortMerge.maxFilesPerPartition readers are open at once regardless of this limit -- " + + "use that config to cap memory on a sorted table.") .intConf .checkValue(v => v > 0, "Data file concurrency limit must be positive") .createWithDefault(1) diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala index 96626631393..fdd38ef1993 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala @@ -27,7 +27,9 @@ import scala.util.control.NonFatal import org.apache.hadoop.conf.Configuration import org.apache.spark.internal.Logging import org.apache.spark.sql.SparkSession +import org.apache.spark.sql.catalyst.expressions.SortOrder +import org.apache.comet.serde.ExprOuterClass.Expr import org.apache.comet.util.ClassLoaders /** @@ -1158,26 +1160,56 @@ object IcebergReflection extends Logging { * types with either layout (e.g. geometry). */ def pageIndexUnsupportedColumns(schema: Any): Set[String] = { + topLevelColumnTypes(schema, "page-index-unsupported columns") + .getOrElse(Seq.empty) + .collect { + case (name, typeStr) + if typeStr.startsWith("decimal(") || typeStr == "uuid" || + typeStr.startsWith("fixed[") || typeStr == "binary" => + name + } + .toSet + } + + /** + * Names of top-level columns whose Iceberg sort order can differ from Spark's comparison of the + * Spark type they map to, so a native k-way merge keyed on the Spark value could mis-order + * rows. Today that is UUID: Iceberg maps UUID to Spark StringType (see TypeToSparkType) but + * sorts by its own UUID comparator, not by the canonical string, so the file order and a string + * comparison can disagree. reportableOrdering refuses a sort key in this set. v1 only reports + * identity, top-level sort keys, so top-level columns are enough. + * + * Returns None if the schema could not be read. That is the fail-closed answer: we cannot rule + * out a UUID (or any other unsafe) sort key, so the caller must refuse the ordering rather than + * assume it is safe -- reporting an ordering the native merge cannot reproduce returns silently + * mis-ordered rows, because Spark has already dropped the Sort by then. + */ + def orderingUnsafeColumns(schema: Any): Option[Set[String]] = { + topLevelColumnTypes(schema, "ordering-unsafe columns").map { columns => + columns.collect { case (name, typeStr) if typeStr == "uuid" => name }.toSet + } + } + + /** + * Read the top-level columns of an Iceberg schema by reflection, yielding (name, typeStr) per + * column. Returns None if the schema could not be read, letting each caller pick its own + * fail-open vs fail-closed answer. + */ + private def topLevelColumnTypes(schema: Any, context: String): Option[Seq[(String, String)]] = { import scala.jdk.CollectionConverters._ try { val columns = getMethod(schema.getClass, "columns") .invoke(schema) .asInstanceOf[java.util.List[_]] - columns.asScala.flatMap { column => + Some(columns.asScala.iterator.map { column => val name = getMethod(column.getClass, "name").invoke(column).asInstanceOf[String] val typeStr = getMethod(column.getClass, "type").invoke(column).toString - if (typeStr.startsWith("decimal(") || typeStr == "uuid" || typeStr.startsWith("fixed[") || - typeStr == "binary") { - Some(name) - } else { - None - } - }.toSet + (name, typeStr) + }.toSeq) } catch { case e: Exception => - logWarning( - s"Failed to inspect schema for page-index-unsupported columns: ${e.getMessage}") - Set.empty[String] + logWarning(s"Failed to inspect schema for $context: ${e.getMessage}") + None } } @@ -2362,7 +2394,18 @@ case class CometIcebergNativeScanMetadata( globalFieldIdMapping: Map[String, Int], catalogProperties: Map[String, String], catalogName: Option[String], - fileFormat: String) + fileFormat: String, + // The ordering the scan will report to Spark, evaluated once in CometScanRule and carried here + // so CometIcebergNativeScanExec.outputOrdering and the proto serde read the same value instead + // of re-running the gate. Empty means "report no ordering". Defaulted so the extract() path, + // which runs before output/ordering are known, can build the metadata and CometScanRule fills + // it in via copy(). + reportedOrdering: Seq[SortOrder] = Nil, + // The same ordering already bound to proto at planning time (against the scan's output), so the + // executor-side serde writes it directly without re-binding. Bound once in CometScanRule; if + // the binding fails there the scan falls back to Spark rather than converting and failing at + // task start. Index-aligned with reportedOrdering. + reportedOrderingProto: Seq[Expr] = Nil) object CometIcebergNativeScanMetadata extends Logging { diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index c7c9a31909d..303ecc7f53c 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -32,7 +32,7 @@ import scala.jdk.CollectionConverters._ import org.apache.hadoop.conf.Configuration import org.apache.spark.internal.Logging import org.apache.spark.sql.SparkSession -import org.apache.spark.sql.catalyst.expressions.{Attribute, DynamicPruningExpression, Expression, GenericInternalRow, InputFileBlockLength, InputFileBlockStart, InputFileName} +import org.apache.spark.sql.catalyst.expressions.{Attribute, DynamicPruningExpression, Expression, GenericInternalRow, InputFileBlockLength, InputFileBlockStart, InputFileName, SortOrder} import org.apache.spark.sql.catalyst.rules.Rule import org.apache.spark.sql.catalyst.util.{sideBySide, ArrayBasedMapData, DateTimeUtils, GenericArrayData, MetadataColumnHelper} import org.apache.spark.sql.catalyst.util.ResolveDefaultColumns.getExistenceDefaultValues @@ -50,6 +50,7 @@ import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, isSpark35Plu import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection, IcebergStorageSchemes} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported, readFieldId} +import org.apache.comet.serde.ExprOuterClass.Expr import org.apache.comet.serde.operator.{CometIcebergNativeScan, CometNativeScan} import org.apache.comet.shims.{CometTypeShim, ShimCometStreaming, ShimFileFormat, ShimSubqueryBroadcast} @@ -1039,16 +1040,68 @@ 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) + // reportedOrdering: the sort prefix Comet will advertise. stayNative: whether Comet may + // keep the scan native even when that prefix is shorter than Iceberg's full reported + // order -- + // true when the whole order is honoured, or when the first field Comet cannot report is an + // unprojected attribute (nothing above can require an ordering that reaches it, so Spark + // dropped no Sort and an unordered read of the tail is correct). It is false when a + // projected UUID / transform blocks (Spark may already have dropped a Sort on it) or when + // the schema could not be read (we cannot rule out an unsafe sort key). + val (reportedOrdering: Seq[SortOrder], stayNative: Boolean) = { + if (!icebergReportsOrdering) { + (Nil, true) + } else { + IcebergReflection.orderingUnsafeColumns(metadata.tableSchema) match { + case Some(unsafe) => + val decision = CometIcebergNativeScan + .orderingDecision(scanExec.ordering, scanExec.output, unsafe) + (decision.reported, decision.stayNative) + case None => (Nil, false) + } + } + } + // Bind the reported prefix to proto now, against the same output the gate used, so the + // executor-side serde writes it directly and a binding drift becomes a planning-time + // fallback rather than a task-start failure. serializeReportedOrdering returns Some(Nil) + // for an empty prefix (the stay-native unordered read), so an empty prefix does not force a + // fallback on its own. + val reportedOrderingProto: Option[Seq[Expr]] = + CometIcebergNativeScan.serializeReportedOrdering(reportedOrdering, scanExec.output) + val orderingHonored: Boolean = { + val honored = stayNative && reportedOrderingProto.isDefined + if (!honored) { + fallbackReasons += "Iceberg reports a sort order the native scan cannot guarantee " + + "(a transform or unsafe-type sort key that a required ordering could reach, the " + + "schema could not be read, or the order could not be serialized); staying on Spark " + + "so the reported ordering is preserved" + } + honored + } + if (schemaSupported && fileIOCompatible && formatVersionSupported && defaultValuesSupported && schemaTypesSupported && encryptionKeyLengthSupported && taskValidation.allParquet && allSupportedFilesystems && allLocationsOpenable && metadataSchemeSupported && partitionTypesSupported && unifiedPartitionTypeSupported && transformFunctionsSupported && deleteFileTypesSupported && dppSubqueriesSupported && - nestedFieldsSupported) { + nestedFieldsSupported && orderingHonored) { CometBatchScanExec( scanExec.clone().asInstanceOf[BatchScanExec], runtimeFilters = scanExec.runtimeFilters, - nativeIcebergScanMetadata = Some(metadata)) + nativeIcebergScanMetadata = Some( + metadata.copy( + reportedOrdering = reportedOrdering, + // Safe: orderingHonored is in the guard above, so when Iceberg reports an order we + // only reach here if the binding succeeded (Some). + reportedOrderingProto = reportedOrderingProto.getOrElse(Nil)))) } else { withFallbackReasons(scanExec, fallbackReasons.toSet) } diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala index 335f02ff034..72f3ca6de70 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala @@ -44,8 +44,9 @@ import org.apache.comet.DataTypeSupport.isComplexType import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.serde.{CometOperatorSerde, OperatorOuterClass} +import org.apache.comet.serde.ExprOuterClass.Expr import org.apache.comet.serde.OperatorOuterClass.{Operator, SparkStructField} -import org.apache.comet.serde.QueryPlanSerde.serializeDataType +import org.apache.comet.serde.QueryPlanSerde.{exprToProto, serializeDataType} object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] with Logging { @@ -944,6 +945,127 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit Some(builder.setIcebergScan(icebergScanBuilder).build()) } + /** + * The longest leading prefix of an Iceberg-reported sort order that the native scan can honour + * (empty if even the first field cannot be reported). The single caller is CometScanRule, which + * pairs this with orderingDecision to decide fallback-vs-stay-native, then stashes the result + * on the scan metadata; CometIcebergNativeScanExec.outputOrdering (what tells Spark the scan is + * sorted) and the proto serialization both read that one stashed value, so 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 the prefix stops + * at them. Checking exprToProto here, not just in the proto path, keeps the advertised order + * and the serialized order identical. + * + * 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. On Spark + * 3.4/3.5/4.0 outputOrdering blanks when multiple InputPartitions share a partition key + * (DataSourceV2ScanExecBase filters `ordering` on `groupedParts.forall(_.parts.length <= 1)`), + * NOT when a partition holds multiple files (note: Spark master no longer blanks at all, so + * this is version-specific). Iceberg's SortOrderAnalyzer.hasUniquePartitionKeys refuses to + * report in exactly that same case, so on the Iceberg path scanExec.ordering and + * scanExec.outputOrdering are equal and the many-files-in-one-partition case this merge targets + * is one Spark does not blank. Reading `ordering` therefore matches outputOrdering today; it + * leans on that Iceberg invariant to stay safe -- if a future Iceberg relaxes + * hasUniquePartitionKeys, switch this read to scanExec.outputOrdering on the versions that + * blank it. + * + * `unsafeColumns` are top-level columns whose Iceberg sort order can differ from Spark's + * comparison of the mapped Spark type (today: UUID, which Iceberg maps to StringType but sorts + * by its own comparator -- see IcebergReflection.orderingUnsafeColumns). A sort key on such a + * column is refused, because the files are sorted by an order the native string merge would not + * reproduce. This cannot be detected from the Spark type alone (UUID looks like a plain + * string), so the caller passes the names in. + */ + def reportableOrdering( + ordering: Option[Seq[SortOrder]], + output: Seq[Attribute], + unsafeColumns: Set[String] = Set.empty): Seq[SortOrder] = { + // The longest leading run of reportable fields. Reporting a prefix is safe because Spark's + // SortOrder.orderingSatisfies matches positionally from index 0, so any required ordering above + // the scan can reach at most this prefix. Stopping at the first non-reportable field is right; + // whether Comet may still stay native there (vs fall back) is decided by orderingDecision. + ordering match { + case Some(orders) => orders.takeWhile(isReportable(_, output, unsafeColumns)) + case None => Nil + } + } + + /** + * The reportability gate's decision for CometScanRule: the sort prefix Comet will advertise + * (`reported`), plus whether Comet may stay on the native scan (`stayNative`) even when that + * prefix is shorter than Iceberg's full reported order. + * + * Comet may stay native iff the first field it cannot report is an unprojected attribute: + * nothing above the scan can require an ordering that reaches such a field (it is not in the + * output, and orderingSatisfies matches positionally from index 0), so Spark has dropped no + * Sort that depends on it, and an unordered read of the tail is correct -- Comet advertises + * just the safe prefix. A projected-but-unsafe field (a UUID, whose Iceberg byte order differs + * from Spark's string comparison) or a non-column transform can appear in a required ordering + * that Spark may already have eliminated a Sort on, so there Comet must fall back to Spark's + * own reader. + */ + case class OrderingDecision(reported: Seq[SortOrder], stayNative: Boolean) + + def orderingDecision( + ordering: Option[Seq[SortOrder]], + output: Seq[Attribute], + unsafeColumns: Set[String]): OrderingDecision = { + val reported = reportableOrdering(ordering, output, unsafeColumns) + val full = ordering.getOrElse(Nil) + val stayNative = + reported.length == full.length || // Iceberg reported nothing, or we honour the whole order + isUnprojectedAttribute(full(reported.length), output) // blocker is unprojected -> safe + OrderingDecision(reported, stayNative) + } + + private def isUnprojectedAttribute(order: SortOrder, output: Seq[Attribute]): Boolean = + order.child match { + case a: AttributeReference => !output.exists(_.exprId == a.exprId) + case _ => false + } + + private def isReportable( + order: SortOrder, + output: Seq[Attribute], + unsafeColumns: Set[String]): Boolean = + isIdentityProjected(order, output) && !isUnsafeColumn(order, unsafeColumns) && + exprToProto(order, output).isDefined + + private def isUnsafeColumn(order: SortOrder, unsafeColumns: Set[String]): Boolean = + order.child match { + case a: AttributeReference => unsafeColumns.contains(a.name) + case _ => false + } + + private def isIdentityProjected(order: SortOrder, output: Seq[Attribute]): Boolean = + order.child match { + case a: AttributeReference => output.exists(_.exprId == a.exprId) + case _ => false + } + + /** + * Binds the reported ordering to proto once, at planning time, against the same `output` the + * gate used. Returns None if any SortOrder fails to serialize, so CometScanRule can fall back + * to Spark cleanly instead of converting and then failing at task start. reportableOrdering + * already checks each order serializes against this output, so None here means the binding + * drifted -- treat the ordering as unreportable rather than raising at execution time. + */ + def serializeReportedOrdering( + reportedOrdering: Seq[SortOrder], + output: Seq[Attribute]): Option[Seq[Expr]] = { + if (reportedOrdering.isEmpty) { + Some(Nil) + } else { + val protoOrders = reportedOrdering.map(exprToProto(_, output)) + if (protoOrders.forall(_.isDefined)) Some(protoOrders.map(_.get)) else None + } + } + /** * Serializes partitions from inputRDD at execution time. * @@ -1033,6 +1155,16 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit commonBuilder.setMetadataLocation(metadata.metadataLocation) commonBuilder.setDataFileConcurrencyLimit( CometConf.COMET_ICEBERG_DATA_FILE_CONCURRENCY_LIMIT.get()) + // sortMerge.enabled = false keeps the scan native and still honours the reported order, but + // via the spillable SortExec rather than the k-way merge. Express that as "merge at most 0 + // files per partition" so the planner always takes the sort path; when enabled, carry the + // configured cap. Lives on the proto next to data_file_concurrency_limit, so the default is + // defined once (in CometConf) and the native side reads common.max_files_per_partition. + commonBuilder.setMaxFilesPerPartition(if (CometConf.COMET_ICEBERG_SORT_MERGE_ENABLED.get()) { + CometConf.COMET_ICEBERG_SORT_MERGE_MAX_FILES_PER_PARTITION.get() + } else { + 0 + }) metadata.catalogName.foreach(commonBuilder.setCatalogName) (metadata.catalogProperties + ioTimeoutProperty()).foreach { case (key, value) => commonBuilder.putCatalogProperties(key, value) @@ -1047,6 +1179,13 @@ object CometIcebergNativeScan extends CometOperatorSerde[CometBatchScanExec] wit commonBuilder.addRequiredSchema(field.build()) } + // The reported ordering was bound to proto at planning time (metadata.reportedOrderingProto) + // against this same output, so a binding failure already fell the scan back to Spark in + // CometScanRule -- here we just write the pre-bound protos, no re-binding or exec-time throw. + if (metadata.reportedOrderingProto.nonEmpty) { + commonBuilder.addAllTableSortOrders(metadata.reportedOrderingProto.asJava) + } + // Load Iceberg classes once (avoid repeated class loading in loop) val contentScanTaskClass = IcebergReflection.loadClass(IcebergReflection.ClassNames.CONTENT_SCAN_TASK) diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergNativeScanExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergNativeScanExec.scala index e058a5c0916..65fc0aa3022 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergNativeScanExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergNativeScanExec.scala @@ -128,9 +128,25 @@ case class CometIcebergNativeScanExec( // Only accessed during execution, not planning def numPartitions: Int = perPartitionData.length + // Spark's EnsureRequirements runs before Comet converts this scan, so it already eliminates a + // redundant shuffle based on the vanilla BatchScanExec's KeyGroupedPartitioning (storage- + // partitioned join) while the leaf is still a BatchScanExec. This node therefore does not need + // to re-report partitioning: UnknownPartitioning is enough, and it keeps the whole SPJ + // negotiation (including commonPartitionValues push-down under AQE) on BatchScanExec, which is + // the only node Spark can push into. override lazy val outputPartitioning: Partitioning = UnknownPartitioning(numPartitions) - override lazy val outputOrdering: Seq[SortOrder] = Nil + // Report the Iceberg-reported ordering (gated to what the native per-partition merge honours) so + // Spark elides redundant sorts above the scan. The gate ran once in CometScanRule and stashed the + // result on nativeIcebergScanMetadata, so this and the proto serialization advertise the exact + // same order. Nil on canonicalized instances (originalPlan / metadata nulled) -- they are never + // executed and their ordering is irrelevant to equality. + override lazy val outputOrdering: Seq[SortOrder] = + if (originalPlan == null || nativeIcebergScanMetadata == null) { + Nil + } else { + nativeIcebergScanMetadata.reportedOrdering + } /** * Maps Iceberg V2 custom metric types to standard Spark metric types for better UI formatting. diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergSortMergeReadSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergSortMergeReadSuite.scala new file mode 100644 index 00000000000..1d605cf4f1e --- /dev/null +++ b/spark/src/test/scala/org/apache/comet/CometIcebergSortMergeReadSuite.scala @@ -0,0 +1,960 @@ +/* + * 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)() + private val gateC = AttributeReference("c", 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: reports the longest identity-projected prefix") { + // a is projected and reportable; the transform on b stops the prefix, so only [a] is reported. + val ordering = Seq(SortOrder(gateA, Ascending), SortOrder(Add(gateB, Literal(1)), Ascending)) + assert( + CometIcebergNativeScan.reportableOrdering(Some(ordering), Seq(gateA, gateB)) === + Seq(SortOrder(gateA, Ascending))) + } + + test("gate: a leading column outside the projection reports nothing") { + val ordering = Seq(SortOrder(gateA, Ascending)) + assert(CometIcebergNativeScan.reportableOrdering(Some(ordering), Seq(gateB)).isEmpty) + } + + test("gate: a leading transform (non-AttributeReference) sort child reports nothing") { + val ordering = Seq(SortOrder(Add(gateA, Literal(1)), Ascending)) + assert(CometIcebergNativeScan.reportableOrdering(Some(ordering), Seq(gateA)).isEmpty) + } + + test("gate: absent or empty ordering reports nothing") { + 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) reports nothing") { + // "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) + } + + // orderingDecision decides fallback-vs-stay-native, which reportableOrdering (the advertised + // prefix) alone does not capture: an empty prefix can mean "stay native, read unordered" (safe) + // or "fall back to Spark" (a projected UUID/transform Spark may have dropped a Sort on). + + test("decision: a leading unprojected sort key stays native with no reported ordering") { + // b (the sort key) is not selected. Nothing above the scan can require an ordering that reaches + // it, so Spark dropped no Sort -- Comet stays native and simply reports nothing. + val ordering = Seq(SortOrder(gateB, Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA), Set.empty) + assert(d.reported.isEmpty && d.stayNative) + } + + test("decision: a projected prefix then an unprojected key stays native with the prefix") { + // a is projected + reportable; b is not selected. Report [a] and stay native. + val ordering = Seq(SortOrder(gateA, Ascending), SortOrder(gateB, Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA), Set.empty) + assert(d.reported === Seq(SortOrder(gateA, Ascending)) && d.stayNative) + } + + test("decision: a two-field projected prefix then an unprojected key stays native") { + // a and b are projected + reportable; c is not selected. Report [a, b] and stay native. + val ordering = + Seq(SortOrder(gateA, Ascending), SortOrder(gateB, Ascending), SortOrder(gateC, Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA, gateB), Set.empty) + assert(d.reported === Seq(SortOrder(gateA, Ascending), SortOrder(gateB, Ascending))) + assert(d.stayNative) + } + + test("decision: an unprojected blocker before a projected UUID still stays native") { + // First blocker wins: b (unprojected) stops the prefix at [a], so we stay native even though a + // later field (c) is a projected unsafe/UUID column -- nothing above can require past b. + val ordering = + Seq(SortOrder(gateA, Ascending), SortOrder(gateB, Ascending), SortOrder(gateC, Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA, gateC), Set("c")) + assert(d.reported === Seq(SortOrder(gateA, Ascending)) && d.stayNative) + } + + test("decision: a projected prefix then a projected UUID falls back to Spark") { + // a is projected + reportable; b is a projected UUID (unsafe). The prefix is [a], but the + // blocker b is IN the output, so a required ordering could reach it and Spark may already have + // dropped a Sort on it -- Comet must fall back even though a non-empty prefix is reportable. + // This pins the `stayNative && ...` gate against a regression to the old `reported.nonEmpty`. + val ordering = Seq(SortOrder(gateA, Ascending), SortOrder(gateB, Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA, gateB), Set("b")) + assert(d.reported === Seq(SortOrder(gateA, Ascending)) && !d.stayNative) + } + + test("decision: a leading UUID sort key falls back to Spark") { + // "a" is a projected UUID: a required ordering can reach it and Spark may already have dropped + // a Sort on it, so Comet must fall back rather than read unordered. + val ordering = Seq(SortOrder(gateA, Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA, gateB), Set("a")) + assert(d.reported.isEmpty && !d.stayNative) + } + + test("decision: a leading transform sort key falls back to Spark") { + val ordering = Seq(SortOrder(Add(gateA, Literal(1)), Ascending)) + val d = CometIcebergNativeScan.orderingDecision(Some(ordering), Seq(gateA), Set.empty) + assert(d.reported.isEmpty && !d.stayNative) + } + + test("decision: no reported ordering stays native") { + assert(CometIcebergNativeScan.orderingDecision(None, Seq(gateA), Set.empty).stayNative) + } + + // --------------------------------------------------------------------------------------------- + // 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("unprojected sort key stays native with no reported ordering") { + 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) // c1 is the sort key, not selected below + insertBatches(cat, "t", "(1,'a','P1'),(3,'c','P1')", "(2,'b','P1'),(4,'d','P1')") + + // c1 (the leading sort key) is not in the projection. Nothing above the scan can require an + // ordering that reaches it (orderingSatisfies matches positionally from index 0), so Spark + // drops no Sort -- Comet stays on the native scan and simply reports no ordering, rather than + // falling back to Spark. Holds on any Iceberg build: a non-reporting build reports no ordering + // either way. Result stays correct. + val (_, plan) = checkSparkAnswer(s"SELECT c2 FROM $cat.db.t WHERE c3 = 'P1'") + assert( + nativeScans(plan).nonEmpty, + s"an unprojected sort key must keep the scan native, not fall back to Spark:\n$plan") + nativeScans(plan).foreach { scan => + assert( + scan.outputOrdering.isEmpty, + s"ordering must not be reported when the sort key is not projected:\n$plan") + } + } + } + + test("results are identical with the sort-merge feature on and off") { + 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')", + "(5,'e','P1'),(6,'f','P1')") + val query = s"SELECT id, data FROM $cat.db.t ORDER BY id" + + val enabled = spark.sql(query).collect().toSeq + // withSQLConf returns the block value on Spark 4.x but Unit on 3.4/3.5, so capture the rows + // in a var rather than relying on its return value. + var disabled: Seq[org.apache.spark.sql.Row] = Seq.empty + withSQLConf(CometConf.COMET_ICEBERG_SORT_MERGE_ENABLED.key -> "false") { + disabled = spark.sql(query).collect().toSeq + } + assert(enabled == disabled, "reporting the ordering must not change the result") + } + } + + // ------------------------------------------------------------------------------------------- + // Fallback to Spark. These are the v1 non-targets from + // https://github.com/apache/datafusion-comet/issues/5323: cases where Iceberg reports a sort + // order but the native scan cannot honour it. When Iceberg reports an ordering, EnsureRequirements + // may already have dropped the Sort above the scan, so a native unordered read would be silently + // wrong -- Comet must instead leave the scan on Spark's Iceberg reader. Correctness is checked + // unconditionally (checkSparkAnswer, holds on any Iceberg build); the "stayed on Spark" assertion + // is gated on a reporting build (assertFellBackToSpark). + // + // Triggers for the Spark-fallback path (Comet declines a reported ordering it cannot honour): a + // transform (bucket) sort key (#5339), which Comet's reportableOrdering rejects as a non-column + // sort child. The transform test is gated on Spark 4.0+: Spark converts a transform sort ordering + // only where V2ScanPartitioningAndOrdering threads the function catalog into the scan's + // outputOrdering (4.0+ does, 3.4 does not -- there a transform ordering throws + // _LEGACY_ERROR_TEMP_3054 in Spark before Comet is reached), and it also needs a reporting Iceberg + // build (assertFellBackToSpark gates on that). UUID sort keys are covered by the "ordering-unsafe + // column" gate test since a UUID column has no Spark DDL type. Note that sort-merge disabled is + // NOT a fallback trigger: the scan stays native and honours the order via a spillable sort (see + // the "sort-merge disabled keeps the scan native" test above). + // ------------------------------------------------------------------------------------------- + + test("sort-merge disabled but Iceberg reports an ordering stays native over an SMJ") { + withSortedTables(spjConf ++ Seq(CometConf.COMET_ICEBERG_SORT_MERGE_ENABLED.key -> "false"))( + "a", + "b") { cat => + Seq("a", "b").foreach { t => + spark.sql( + s"CREATE TABLE $cat.db.$t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + replaceSortOrder(cat, "db", t, "c1" -> true) + insertBatches(cat, t, "(1,'a','X'),(2,'b','X')", "(1,'c','X'),(3,'d','X')") + } + val query = s"SELECT a.c1, b.c2 FROM $cat.db.a a JOIN $cat.db.b b ON a.c1 = b.c1" + val (_, plan) = checkSparkAnswer(query) + // Disabling the merge must not push the join's scans back to Spark. + assert( + nativeScans(plan).nonEmpty, + s"sort-merge disabled must not force the SMJ scans back to Spark:\n$plan") + } + } + + test("sort-merge disabled but Iceberg reports an ordering stays native over a group-by") { + withSortedTables(spjConf ++ Seq(CometConf.COMET_ICEBERG_SORT_MERGE_ENABLED.key -> "false"))( + "t") { cat => + spark.sql( + s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + replaceSortOrder(cat, "db", "t", "c1" -> true) + insertBatches(cat, "t", "(1,'a','X'),(2,'b','X')", "(1,'c','X'),(3,'d','X')") + val query = s"SELECT c1, COUNT(*) FROM $cat.db.t GROUP BY c1" + val (_, plan) = checkSparkAnswer(query) + assert( + nativeScans(plan).nonEmpty, + s"sort-merge disabled must not force the group-by scan back to Spark:\n$plan") + } + } + + test("fallback: a transform (bucket) sort order stays on Spark (#5339)") { + // Spark converts a transform sort ordering only on 4.0+ (see the section note); on 3.4 the + // query would throw in Spark before Comet is reached, so skip there. + assume(isSpark40Plus, "transform sort orderings are only convertible on Spark 4.0+") + withSortedTables(spjConf)("t") { cat => + spark.sql( + s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + // Sort by bucket(8, c1): Iceberg reports the ordering as a bucket transform. Spark 4.0+ + // converts it, but Comet's reportableOrdering rejects the non-AttributeReference sort child, + // so Comet stays on Spark (v1 identity-only, #5339). + replaceSortOrderBucket(cat, "db", "t", "c1", 8) + insertBatches(cat, "t", "(1,'a','X'),(2,'b','X')", "(3,'c','X'),(4,'d','X')") + val query = s"SELECT c1, c2 FROM $cat.db.t" + val (_, plan) = checkSparkAnswer(query) + assertFellBackToSpark(query, plan) + } + } + + // Sort-merge join: sorted + storage-partitioned inputs let EnsureRequirements drop both the + // required-child Sort AND the Exchange. Target: zero of each (assume-gated on the reporting). + + test("sort-merge join over co-partitioned sorted tables drops the sort and the shuffle") { + withSortedTables(spjConf)("a", "b") { cat => + Seq("a", "b").foreach { t => + spark.sql( + s"CREATE TABLE $cat.db.$t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + replaceSortOrder(cat, "db", t, "c1" -> true) + insertBatches( + cat, + t, + s"(1,'${t}1','X'),(2,'${t}2','X')", + s"(3,'${t}3','X'),(4,'${t}4','X')") + } + + val query = + s"SELECT t1.c1, t1.c2, t2.c2 FROM $cat.db.a t1 JOIN $cat.db.b t2 ON t1.c1 = t2.c1" + val (_, plan) = checkSparkAnswer(query) + + // Shuffle elimination needs only the grouping (upstream Iceberg), so assert it always. + assert(countShuffles(plan) == 0, s"storage-partitioned join must not shuffle:\n$plan") + // Sort elimination needs the reported ordering (fork Iceberg today), so gate that one. + assumeOrderingReported(plan) + assert(countSorts(plan) == 0, s"reported ordering must eliminate the join sort:\n$plan") + } + } + + test("storage-partitioned join on the bucket key drops the shuffle") { + withSortedTables(spjConf)("a", "b") { cat => + // No sort order set: this isolates the SPJ shuffle-elimination (KeyGroupedPartitioning), + // independent of the ordering path. A join sort may remain; the Exchange must not. + Seq("a", "b").foreach { t => + spark.sql( + s"CREATE TABLE $cat.db.$t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + insertBatches( + cat, + t, + s"(1,'${t}1','X'),(2,'${t}2','X')", + s"(3,'${t}3','X'),(4,'${t}4','X')") + } + + val query = + s"SELECT t1.c1, t1.c2, t2.c2 FROM $cat.db.a t1 JOIN $cat.db.b t2 ON t1.c1 = t2.c1" + val (_, plan) = checkSparkAnswer(query) + + // Comet no longer reports partitioning; EnsureRequirements eliminates the shuffle on the + // BatchScanExec before Comet converts the scan, so this holds on any Iceberg with grouping. + assert(countShuffles(plan) == 0, s"storage-partitioned join must not shuffle:\n$plan") + } + } + + // AQE runtime partition-value pushdown. Under AQE Spark computes the exact common set of + // partition values once both sides materialize and pushes it into the scan + // (EnsureRequirements.populateCommonPartitionInfo, which matches only `case scan: BatchScanExec`). + // Because Comet converts the scan AFTER EnsureRequirements, that push-down lands on the vanilla + // BatchScanExec, and Comet's execution partitions come from originalPlan.inputRDD -- already + // padded to the merged set. So the native leaf stays aligned with Spark's spec without Comet + // reporting any partitioning itself. Offset partition-value sets (a: c1 in 1..4, b: 3..6) make + // the merged/padded set differ from each side's own list, exactly where a misalignment would + // surface. + // + // The asserts are hard, not `assume`d: correctness (Comet vs vanilla Spark) must hold, both sides + // must stay native (a fallback to Spark's BatchScanExec would let correctness pass vacuously), + // and the join must not shuffle (a shuffle would re-partition correctly and mask the pushdown). + // Only Iceberg is a precondition (grouping is upstream, unlike ordering). + + /** a: c1 in 1..4, b: c1 in 3..6, both bucketed the same way -- offset partition-value sets. */ + private def createAqePartitionedPair(cat: String): Unit = { + Seq("a", "b").foreach { t => + spark.sql( + s"CREATE TABLE $cat.db.$t (c1 INT, c2 STRING) USING iceberg " + + "PARTITIONED BY (bucket(8, c1))") + } + insertBatches(cat, "a", "(1,'a1'),(2,'a2')", "(3,'a3'),(4,'a4')") + insertBatches(cat, "b", "(3,'b3'),(4,'b4')", "(5,'b5'),(6,'b6')") + } + + private val aqeSpjConf: Seq[(String, String)] = + spjConf.filterNot(_._1 == "spark.sql.adaptive.enabled") :+ + ("spark.sql.adaptive.enabled" -> "true") + + private def assertNativeSpj(plan: SparkPlan): Unit = { + assert( + nativeScans(plan).length == 2, + s"both join sides must stay on the native Iceberg scan (no fallback):\n$plan") + assert(countShuffles(plan) == 0, s"storage-partitioned join must not shuffle:\n$plan") + } + + // Left outer is the sensitive case: every left row must appear, so if the leaf's partition + // indices are shifted relative to Spark's padded spec, a left row whose key exists on the right + // gets a wrong NULL (or a match against the wrong group) -- caught by checkSparkAnswer. An inner + // join would only drop the same unmatched rows on both sides and can pass despite misalignment. + test("AQE storage-partitioned left outer join with runtime partition pushdown is correct") { + withSortedTables(aqeSpjConf)("a", "b") { cat => + createAqePartitionedPair(cat) + val query = + s"SELECT t1.c1, t1.c2, t2.c2 FROM $cat.db.a t1 LEFT OUTER JOIN $cat.db.b t2 " + + "ON t1.c1 = t2.c1" + val (_, plan) = checkSparkAnswer(query) + assertNativeSpj(plan) + } + } + + test("AQE storage-partitioned inner join with runtime partition pushdown is correct") { + val innerConf = aqeSpjConf :+ + ("spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabled" -> "true") + withSortedTables(innerConf)("a", "b") { cat => + createAqePartitionedPair(cat) + val query = + s"SELECT t1.c1, t1.c2, t2.c2 FROM $cat.db.a t1 JOIN $cat.db.b t2 ON t1.c1 = t2.c1" + val (_, plan) = checkSparkAnswer(query) + assertNativeSpj(plan) + } + } + + // Grouped aggregate / distinct: ReplaceHashWithSortAgg + a satisfied ClusteredDistribution. + + test("group-by on the sort key drops the sort and the shuffle") { + withSortedTables(spjConf)("t") { cat => + spark.sql( + s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + replaceSortOrder(cat, "db", "t", "c1" -> true) + insertBatches(cat, "t", "(1,'a','X'),(2,'b','X')", "(1,'c','X'),(3,'d','X')") + + // No trailing ORDER BY: checkSparkAnswer compares order-insensitively, so sort = 0 is a + // legitimate target (a grouped aggregate over sorted + clustered input needs neither). + val query = s"SELECT c1, COUNT(*) FROM $cat.db.t GROUP BY c1" + val (_, plan) = checkSparkAnswer(query) + + assert(countShuffles(plan) == 0, s"grouped aggregate must not shuffle:\n$plan") + assumeOrderingReported(plan) + assert(countSorts(plan) == 0, s"grouped aggregate over sorted input must not sort:\n$plan") + } + } + + test("distinct on the sort key is correct") { + withSortedTables(spjConf)("t") { cat => + spark.sql( + s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + replaceSortOrder(cat, "db", "t", "c1" -> true) + insertBatches(cat, "t", "(1,'a','X'),(2,'b','X')", "(1,'c','X'),(3,'d','X')") + + checkSparkAnswer(s"SELECT DISTINCT c1 FROM $cat.db.t ORDER BY c1") + } + } + + // Window: WindowExec.requiredChildOrdering (partitionSpec ++ orderSpec) + ClusteredDistribution. + + test("window partitioned + ordered on the sort keys drops the sort and the shuffle") { + withSortedTables(spjConf)("t") { cat => + spark.sql( + s"CREATE TABLE $cat.db.t (c1 INT, c2 STRING, c3 STRING) USING iceberg " + + "PARTITIONED BY (bucket(4, c1))") + replaceSortOrder(cat, "db", "t", "c1" -> true, "c2" -> true) + insertBatches(cat, "t", "(1,'a','X'),(2,'b','X')", "(1,'c','X'),(2,'d','X')") + + 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) + + assert(countShuffles(plan) == 0, s"window over clustered input must not shuffle:\n$plan") + assumeOrderingReported(plan) + assert(countSorts(plan) == 0, s"window over sorted input must not sort:\n$plan") + } + } + + // Limit: TakeOrderedAndProjectExec degenerates to a cheap take when the child is already sorted. + + test("order-by-limit over a sorted table is correct") { + 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')", + "(5,'e','P1'),(6,'f','P1')") + + // TakeOrderedAndProject is planned for ORDER BY ... LIMIT and, when the child is already + // reported sorted, degenerates to a cheap bounded take instead of a full top-N heap. + // Correctness (including the LIMIT cut over the merged order) is the guarantee here. + val (_, plan) = checkSparkAnswer(s"SELECT id, data FROM $cat.db.t ORDER BY id LIMIT 3") + assumeOrderingReported(plan) + } + } +}