anuragmantri opened a new pull request, #55:
URL: https://github.com/apache/datafusion-iceberg/pull/55
## Which issue does this PR close?
- Part 2 of apache/iceberg-rust#3126. Part 1 (per-file sort order on
`FileScanTask`) merged in apache/iceberg-rust#3128 and is included in the
pinned iceberg-rust rev.
## What changes are included in this PR?
`IcebergTableScan` can now report an output ordering to DataFusion when the
data files it reads are already sorted, so DataFusion can drop the `SortExec`
it would otherwise add above the scan for an `ORDER BY` or a sort-merge join.
The feature is off by default. It is enabled per session with `SET
iceberg.planning.preserve_data_ordering = true`, after registering the new
`IcebergDataFusionConfig` extension. The option mirrors Spark's
`spark.sql.iceberg.planning.preserve-data-ordering` (apache/iceberg#16750).
When it is on, `TableProvider::scan` lists the data files up front. The scan
reports an ordering only if:
- there are at most `iceberg.planning.max_merge_files` files (default 64,
matching Comet's `sortMerge.maxFilesPerPartition` in
apache/datafusion-comet#5331), and
- every file records the same `sort_order_id`, and it resolves to a sorted
order.
It reports the leading fields of that order that are identity transforms of
projected top-level columns, stopping at the first float, double, or UUID
field. Arrow orders -0.0 before 0.0 and puts negative NaN first, whereas Spark,
the usual writer of sorted files, treats the zeros as equal and puts NaN last.
The Iceberg spec doesn't define UUID ordering (apache/iceberg#14216). In every
other case the scan is planned exactly as before.
Two files with the same sort order can have overlapping key ranges, so
`execute()` reads each file as its own stream and k-way merges them with
DataFusion's `StreamingMergeBuilder`, which is the machinery
`SortPreservingMergeExec` uses. A pushed-down `LIMIT` becomes the merge's
`fetch`. All the readers are clones of one `ArrowReader`, so they share its
delete-file cache.
The scan keeps a single output partition, so this doesn't depend on
apache/iceberg-rust#2671, and it composes with it. DataFusion's output ordering
is per partition, so once #2671 groups files into partitions, each group can be
merged the same way.
New public API:
- `IcebergDataFusionConfig` and `IcebergPlanningConfig`, the session options.
- `IcebergTableScan::sorted_tasks` and
`IcebergTableScan::with_sorted_tasks`. Since #20, a scan can be rebuilt from
its accessors, for example to ship it to another process. A sorted scan rebuilt
without its listed files would read them unordered beneath a plan that already
dropped the sort, so the rebuild path needs them.
When the option is on but a scan falls back, the files are listed again in
`execute()`, as without the option. The table's manifest cache makes the second
listing cheap.
Merging opens every file at once. The cap bounds that. Bounding it instead
by how many files' key ranges overlap would need per-file column bounds on
`FileScanTask`, which is follow-up work.
## Are these changes tested?
Yes.
- **Unit tests** for mapping a sort order to a DataFusion ordering:
- direction and null order,
- stopping at bucket, truncate, and day transforms,
- stopping at float, double, UUID, nested, unprojected, and missing fields.
- **Integration tests** over tables with hand-written sorted files:
- files with overlapping key ranges and duplicate keys merge into sorted
order, with struct and list columns carried through;
- `LIMIT`;
- no ordering when files disagree, have no sort order id, or have the
unsorted order, and the plan then matches the one without the option;
- the option is off by default, and the cap is honored;
- `ORDER BY` loses its `SortExec`, while `ORDER BY ... DESC` and the
option-off case keep it;
- a sort-merge join with `target_partitions = 4` has no `SortExec`;
- a rebuilt scan keeps its ordering;
- positional deletes are applied on the merge path.
- **Mutation checks.** Replacing the merge with concatenation fails 6 of the
8 integration tests. Dropping the check that every file has the same sort order
fails the disagreement test.
## AI Disclosure
I'm new to Rust. I used Claude Code (Claude Opus 5.5) to write this change
and its tests, and I reviewed it manually.
--
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]