feat(iceberg): thread per-file sort order into FileScanTask - #3128
Conversation
e285cc5 to
187999e
Compare
|
@mbutrovich - Since you have context on this work, could you please take a look? Thanks! |
anoopj
left a comment
There was a problem hiding this comment.
Code looks right to me. Just one comment about unsorted files.
| pub unified_partition_type: Option<Arc<StructType>>, | ||
|
|
||
| /// The sort order that this file's rows are sorted by, resolved from the data file's | ||
| /// `sort_order_id` against the table's known sort orders. `None` if the file has no |
There was a problem hiding this comment.
Note that Iceberg reserves 0 for unsorted. Per the spec:
Order id 0 is reserved for the unsorted order.
So a file with sort_order_id = 0 will resolve to unsorted order whereas a missing or unresolvable id gives None. So the datafusion consumer will need to check for both. e.g. something like sort_order.as_ref().is_some_and(|o| !o.is_unsorted())
May be a good idea to clarify this in the docs.
There was a problem hiding this comment.
Good catch. I changed it so that Some is returned only when the file has valid sort order (>0).
|
Thanks for the review @anoopj. I addressed your feedback. |
|
@CTTY @blackmwk @laskoviymishka Could you please review? |
laskoviymishka
left a comment
There was a problem hiding this comment.
Nice groundwork, and thanks for picking this up! threading the resolved sort order down through the manifest-entry contexts into with_sort_order reads cleanly, and handling the reserved unsorted order (id 0) with !order.is_unsorted() is the right move (that also covers @anoopj's unsorted-files question).
I'd hold this before merging, mostly around how it behaves once Part 2 starts consuming it.
The not_implemented serde guard on the new sort_order field is the main one. Unlike the other guarded fields, SortOrder already derives serde and round-trips in table metadata today — so the guard isn't working around a real limitation, it just means any task with a resolved order errors on serialize. Since populating this field is the whole point of the PR, I'd rather sort the serde story out now than ship a public field that breaks the moment it's actually used.
The other is that we resolve sort_order_id and then drop the raw id, so "never sorted", "reserved unsorted", and "id didn't resolve" all collapse to None. That's fine for reading, but it's exactly the distinction the sorted-scan optimizer will need — a file whose order definition was dropped is still physically sorted. Carrying the raw sort_order_id alongside the resolved order (or at least warning on an unresolvable non-zero id) keeps that recoverable.
A few smaller things I'd like to settle before merge:
- the shared testdata fixture change makes the file-4 assertion pass vacuously and touches four unrelated test modules — a dedicated fixture or an inline metadata tweak would be cleaner
ManifestFileContextnow carries the fullTableMetadataRefwhere a narrow sort-orders map would match the existingunified_partition_typepattern- an aggregate count assertion in the new test, plus a couple of doc/idiom nits, all inline
Once those are addressed, happy to take another pass and approve.
| case_sensitive: bool, | ||
| partition_spec: Option<PartitionSpecRef>, | ||
| unified_partition_type: Option<Arc<StructType>>, | ||
| table_metadata: TableMetadataRef, |
There was a problem hiding this comment.
Carrying the whole TableMetadataRef into every ManifestFileContext works, but it hands each context a much wider capability than it needs — the existing unified_partition_type pattern precomputes the derived value in PlanContext and carries only that.
Sort-order resolution is genuinely per-entry, so I get why full metadata is tempting. But we could build an Arc<HashMap<i64, SortOrderRef>> once in create_manifest_file_context and carry that instead, which keeps the context narrow and matches the surrounding pattern. wdyt?
There was a problem hiding this comment.
Done. PlanContext now precomputes sort_orders: Arc<HashMap<i64, SortOrderRef>> once in TableScanBuilder::build(), and ManifestFileContext/create_manifest_file_context carry only that map instead of the full TableMetadataRef, matching the unified_partition_type pattern.
| let sort_order = manifest_entry | ||
| .data_file() | ||
| .sort_order_id() | ||
| .and_then(|id| table_metadata.sort_order_by_id(id as i64)) |
There was a problem hiding this comment.
This resolves and then discards the raw sort_order_id, which collapses three different situations into a single None: no id at all, the reserved unsorted order, and an id that doesn't resolve against the table's sort orders.
For reading that's spec-compliant, but for the Part 2 sorted-scan work it's the difference between "this file was never sorted" and "this file is physically sorted but its order definition was dropped" — both arrive as None and the optimizer can't tell them apart.
I'd carry the raw sort_order_id: Option<i32> on FileScanTask alongside the resolved sort_order so a caller can still distinguish those cases. At minimum, a tracing::warn! when a non-zero id fails to resolve would keep this from silently degrading with no signal. wdyt?
There was a problem hiding this comment.
Added FileScanTask::sort_order_id: Option<i32> alongside sort_order, carried through unresolved from DataFile::sort_order_id(). So now sort_order_id distinguishes all three cases (absent, unresolvable, resolved) while sort_order stays None unless the id resolves to a genuinely sorted order. That should give the Part 2 optimizer what it needs even when the order definition itself is gone.
| .data_file() | ||
| .sort_order_id() | ||
| .and_then(|id| table_metadata.sort_order_by_id(id as i64)) | ||
| .filter(|order| !order.is_unsorted()) |
There was a problem hiding this comment.
Small thing: !order.is_unsorted() keys off fields.is_empty() rather than the order id, so a malformed or forward-compat order with a non-zero id and empty fields would also be treated as unsorted. In a valid table that's equivalent to order_id == 0 and it matches Java's isUnsorted(), so it's fine as-is.
The field doc over in task.rs says it resolves "to the reserved unsorted order (id 0, per the spec)", which is a bit more specific than what the check actually does — worth either softening the wording or comparing against SortOrder::UNSORTED_ORDER_ID if we want the doc to be literally true.
There was a problem hiding this comment.
Went with softening the doc rather than switching to UNSORTED_ORDER_ID, since checking is_unsorted() on the resolved order is the behavior I actually want (matches Java's isSorted() gate). Reworded so the doc leads with "an order with no sort fields" and mentions id 0 only as the spec-defined example, rather than implying the check keys off the id.
| current_snapshot.sequence_number(), | ||
| ); | ||
| manifest_list_write | ||
| .add_manifests(vec![data_file_manifest].into_iter()) |
There was a problem hiding this comment.
std::iter::once(data_file_manifest) avoids the throwaway heap Vec here — the sibling helpers use vec![...] because they genuinely have multiple manifests, but this one only ever has the one.
There was a problem hiding this comment.
Fixed for the helper this PR touches. Left the other manifest-list helpers as vec![...] since they're pre-existing code outside this diff.
| .await | ||
| .unwrap(); | ||
|
|
||
| assert_eq!(tasks.len(), 4, "expected all four FileScanTasks"); |
There was a problem hiding this comment.
The per-file assertions are thorough, but nothing asserts the aggregate — a regression that resolved every entry to id 3, or dropped resolution entirely, would still pass three of the four checks.
A single assert_eq!(tasks.iter().filter(|t| t.sort_order.is_some()).count(), 1) up top would catch that class of systemic regression cheaply.
There was a problem hiding this comment.
Added two: an assertion that exactly one task resolves to sort_order.is_some(), and one that three tasks carry a raw sort_order_id. Both would catch the systemic-regression case you described.
| /// unsorted order (id 0, per the spec). | ||
| /// | ||
| /// Note: this reflects only the file's own recorded sort order id, not necessarily the | ||
| /// table's current default sort order — see [`crate::spec::DataFile::sort_order_id`]. |
There was a problem hiding this comment.
The intra-doc link crate::spec::DataFile::sort_order_id points at a method but omits the (), so rustdoc may resolve it as a field or emit a broken_intra_doc_links warning — writing it as sort_order_id() fixes it.
There was a problem hiding this comment.
Fixed, now links to sort_order_id(). Also reran cargo doc with broken_intra_doc_links denied to confirm.
| /// Serde: not yet implemented. | ||
| #[serde(default)] | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| #[serde(serialize_with = "serialize_not_implemented")] |
There was a problem hiding this comment.
This is the one I'd really want to settle before merge. Unlike partition / partition_spec / unified_partition_type, SortOrder already derives Serialize/Deserialize and gets round-tripped in table metadata JSON today — so there's no technical reason for the not_implemented guard here.
The effect is that any FileScanTask whose sort_order is Some(_) errors on serialize. skip_serializing_if covers the None case, but the moment Part 2 populates this field — which is the whole point of the PR — every task with a resolved order silently fails serde on a public field.
I'd drop the serialize_with/deserialize_with guards and use plain #[serde(default, skip_serializing_if = "Option::is_none")], which just works for a round-trippable SortOrder. If there's a wire-format reason to omit it, I'd love a comment spelling that out specifically rather than the not_implemented stub. wdyt?
There was a problem hiding this comment.
Agreed, dropped the not_implemented guard. SortOrder already derives Serialize/Deserialize and the workspace's serde already has the rc feature on, so Option<SortOrderRef> round-trips with plain #[serde(default, skip_serializing_if = "Option::is_none")]. Verified with a doc build too, no wire-format reason to keep it stubbed.
| "default-sort-order-id": 3, | ||
| "sort-orders": [ | ||
| { | ||
| "order-id": 0, |
There was a problem hiding this comment.
Adding order-id 0 here works, but this fixture is shared by four test modules (transaction/append, io/object_cache, the scan tests), and the change flips sort_order_by_id(0) from None to Some(unsorted) for all of them — a silent behavioral change for one test's benefit.
It also makes the file-4 assertion a bit of an illusion: without this entry, sort_order_by_id(0) returns None, and_then short-circuits, and sort_order.is_none() passes either way — the !is_unsorted() filter never actually fires. And the more common Java-written case (id 0 with no explicit order-0 entry) stays untested.
I'd either drop this change and add a dedicated fixture that omits order-id 0 (proving the real short-circuit path), or mutate the metadata clone inline in the helper the way new_unpartitioned does, so the filter branch gets exercised without touching a shared fixture. wdyt?
There was a problem hiding this comment.
Reverted the fixture change. The test now clones the table's metadata and inserts the unsorted order inline (same pattern new_unpartitioned uses to mutate a cloned TableMetadata), so the id-0 case exercises the real !is_unsorted() filter instead of the fixture's absence letting and_then short-circuit to the same result.
|
Thanks for the thorough review, @laskoviymishka. Pushed a commit addressing all of it: |
laskoviymishka
left a comment
There was a problem hiding this comment.
Thanks for iterations! it's good to land now
|
@anuragmantri fix conflicts and it's good to merge |
Resolves each data file's sort_order_id against the table's known sort orders during scan planning and carries the result on FileScanTask as sort_order (Some only when genuinely sorted; the reserved unsorted order, id 0, resolves to None) alongside the raw sort_order_id (kept unresolved so a non-zero id that fails to resolve is still recoverable as "physically sorted"). PlanContext precomputes a narrow sort-orders map once per scan, matching the existing unified_partition_type pattern, rather than carrying the full TableMetadataRef into every ManifestFileContext.
32f31df to
8e8598e
Compare
|
Thanks for the reviews @laskoviymishka and @anoopj |
Which issue does this PR close?
What changes are included in this PR?
Resolves each manifest entry's
DataFile.sort_order_idagainst the table's known sort orders and carries the result on a newFileScanTask.sort_order field. This follows the same PlanContext -> ManifestFileContext -> ManifestEntryContext plumbing already used for unified_partition_type. None means no sort order could be established, either because the file has no recordedsort_order_id, or because the id doesn't resolve against the table's known sort orders.This is groundwork for propagating sort-order awareness into the DataFusion integration (Part 2 of #3126), which needs per-file sort-order data to decide whether a scan's output can be safely reported as sorted.
Are these changes tested?
Yes, a unit test covers all three resolution outcomes:
The manifest-writing setup for this is factored into a new shared fixture helper alongside the existing
setup_*variants.AI Disclosure
I'm new to Rust and this codebase. I used Claude code (Opus 5) to assist me with the PR and the test coverage and I manually reviewed it. Please bear with me as I get familiar with Rust and the codebase.