sunchao commented on code in PR #5331:
URL: https://github.com/apache/datafusion-comet/pull/5331#discussion_r4104363485
##########
native/core/src/execution/planner.rs:
##########
@@ -1844,26 +1845,76 @@ impl PhysicalPlanner {
let metadata_location = common.metadata_location.clone();
let catalog_name = common.catalog_name.clone();
let tasks = parse_file_scan_tasks_from_common(common,
&scan.file_scan_tasks)?;
+ let tasks_len = tasks.len();
let data_file_concurrency_limit =
common.data_file_concurrency_limit as usize;
+ let max_files_per_partition = common.max_files_per_partition
as usize;
- let iceberg_scan = IcebergScanExec::new(
+ // Table sort order Iceberg reported. Empty unless sortMerge
is on and the order
+ // passed the identity gate in CometIcebergNativeScan. The
SortOrder children are
+ // bound references into required_schema, so build the
LexOrdering against it.
+ let ordering: Option<LexOrdering> = if
common.table_sort_orders.is_empty() {
+ None
+ } else {
+ let exprs = common
+ .table_sort_orders
+ .iter()
+ .map(|expr| self.create_sort_expr(expr,
Arc::clone(&required_schema)))
+ .collect::<Result<Vec<PhysicalSortExpr>,
ExecutionError>>()?;
+ LexOrdering::new(exprs)
+ };
+
+ // A per-file-stream k-way merge opens one reader per file at
once. Above the
+ // configured limit we instead read the partition unordered
(bounded by
+ // data_file_concurrency_limit) and sort with a spillable
SortExec, which bounds both
+ // open readers and memory. Both paths still produce sorted
output, so the ordering
+ // Spark eliminated its Sort on is honoured either way. A
limit of 0 (sortMerge
+ // disabled) always takes the sort path.
+ let use_merge = ordering.is_some() && tasks_len <=
max_files_per_partition;
Review Comment:
[P1] Use the full-sort path when a floating-point key precedes another sort
key. With reported ordering `(a DOUBLE ASC, b INT ASC)`, an Iceberg-sorted file
can contain `(-0.0, 2), (+0.0, 1)`: Iceberg's comparator distinguishes the two
zeros. However, `create_sort_expr` normalizes them to the same value, so Spark
requires the row with `b = 1` first. This branch selects
`SortPreservingMergeExec` despite its inputs being unsorted under that
composite comparator, and the merge preserves the incorrect order. This can
change query results: with multiple Spark partitions,
`CometTakeOrderedAndProjectExec` trusts the newly advertised ordering and takes
a local prefix for `ORDER BY a, b LIMIT 2`, potentially discarding the correct
row before the final sort. Route these orderings to the existing spillable
`SortExec` unless Spark-compatible ordering within every file can be
established.
Evidence: A disposable native test wrote two real Parquet files containing
[(-0.0, 2), (+0.0, 1)] and [(-1.0, 0), (1.0, 0)], then constructed an Iceberg
scan protobuf ordered by both columns and executed it through
`PhysicalPlanner::create_plan`. With `max_files_per_partition = 64`, the actual
output was [(-1.0, 0), (-0.0, 2), (+0.0, 1), (1.0, 0)]. With the same files and
cap 0, `SortExec` produced [(-1.0, 0), (+0.0, 1), (-0.0, 2), (1.0, 0)]. The
equality assertion failed reproducibly. Iceberg `Comparators` uses
`Comparator.naturalOrder()` for doubles, while Spark
`SQLOrderingUtil.compareDoubles` treats signed zeros as equal across the
checked versions. Probe source and output remain at
`/tmp/comet-5331-probe/native_test_fragment.rs` and
`/tmp/comet-5331-real-files-probe.log`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]