feat: make the plan nodes inspectable and rebuildable - #20
NoahKusaba wants to merge 5 commits into
Conversation
A plan codec rebuilds a catalog-backed provider from its catalog and table identifier, which needs the constructor. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
mbutrovich
left a comment
There was a problem hiding this comment.
Thanks @NoahKusaba, this is the smaller, codec-shaped change we talked about on apache/iceberg-rust#2862, and the multi-partition commit guard is a good catch. My comments are about whether each new constructor can be called with only what its node exposes, since that's all a decoder in another process has.
A codec decoding these nodes in another process has only their accessors to work from, and two constructors asked for more than that. - IcebergCommitExec::new no longer takes a schema. It was only shown in the verbose display and could not be read back from the node; the display now shows the table's current schema on one line. - IcebergTableScan::new_with_predicate takes the output schema and the projected column names, as schema() and projection() return them, instead of the full table schema and indices, which the scan does not expose. It returns an error when the names are not the schema's fields in order, rather than building a scan whose batches do not match its schema. The integration test now executes a write and commit rebuilt from their accessors, and rebuilds scans from their accessors alone, with and without a projection.
mbutrovich
left a comment
There was a problem hiding this comment.
Thanks for the quick turnaround, @NoahKusaba.
| pub fn new_with_predicate( | ||
| table: Table, | ||
| snapshot_id: Option<i64>, | ||
| schema: ArrowSchemaRef, | ||
| projection: Option<Vec<String>>, | ||
| predicate: Option<Predicate>, | ||
| limit: Option<usize>, | ||
| ) -> Result<Self> { | ||
| // The scan reads the named columns and reports `schema`, so the two | ||
| // must agree, or its batches would not match its schema. | ||
| if let Some(projection) = &projection { | ||
| let fields: Vec<&String> = | ||
| schema.fields().iter().map(|field| field.name()).collect(); | ||
| if !projection.iter().eq(fields.iter().copied()) { | ||
| return plan_err!( | ||
| "IcebergTableScan projection {projection:?} does not match the \ | ||
| fields of its schema {fields:?}" | ||
| ); | ||
| } |
There was a problem hiding this comment.
This follows up on the signature I suggested last round. When projection is Some, it has to equal schema's field names, so it carries nothing that schema doesn't. When it's None, schema has to be the full table schema, and the doc says that isn't checked. On the head commit, passing a one-field projected schema with None still builds a scan whose schema() has one field while its batches have two columns (foo1, foo2).
What do you think about dropping the projection parameter and always selecting schema's field names? For a full schema, that reads the same columns as None. In the iceberg-rust rev this repo uses, select_all and a select of every top-level name resolve to the same field ids in collect_scan_field_ids, and I got identical output from both on the head commit. That removes the unchecked precondition, and the mismatch error and test_scan_rejects_projection_not_matching_schema go away with it.
The trade-off I see is that projection() would return every column name for a scan built without a projection, where it returns None today, and the display would list every column. Is anything relying on that None?
There was a problem hiding this comment.
Agreed, and I checked collect_scan_field_ids at that rev and see the same: None and a select of every top-level name resolve to the same field ids. In 14aaa24, new_with_predicate no longer takes a projection: it always selects schema's field names, so the batches match the schema by construction, and it can't fail. The mismatch check and its test are gone. Nothing relied on projection() being None. It's only read by the display and the tests, and the EXPLAIN tests already show named projections. The display gets clearer too, since a full scan used to print projection:[], the same as an empty projection. projection() keeps its signature and returns the schema's names.
Selecting by name also fixes the same mismatch through the catalog-backed provider, which pairs its cached schema with a freshly loaded table. After a column is added and written to, a scan with no projection used to return that column too. test_scan_after_schema_evolution_reads_provider_columns covers it, and fails on the previous commit.
| let expected = run(scan, &ctx).await?; | ||
| assert_eq!(expected.lines().count(), 5, "one row:\n{expected}"); |
There was a problem hiding this comment.
Comparing each rebuilt plan's output to the original's checks the round trip. The checks on what the originals return are looser: inserted.contains("| 2 |"), contains("alan"), and here a line count of 5, which doesn't show that both columns come back without a projection. The rest of this file asserts results with expect!. Could these assert the full pretty-printed tables with expect![[...]].assert_eq(&expected) instead? Then the test states the rows each plan returns, and the equality with the rebuilt plan covers the rest.
There was a problem hiding this comment.
You're absolutely right, I'll fix that in my next commit.
There was a problem hiding this comment.
Done in 14aaa24: the tests now assert full tables with expect!, and the pinned scan is checked against a later snapshot, so reading the current snapshot instead would show. $snapshots is the exception: its ids, times and paths change every run, so the test checks each row against the snapshot's metadata and compares the summaries in full.
|
Hey @mbutrovich, I just digested your comments, and they make alot of sense to me and I'll look to address them. I also wanted to thank you for all of the time you've put into my many PR's I've made trying to land the Ballista integration, and I appreciate the care and precision you're putting into catching the stray bits of slop that slip past me. I additionally noticed you took the time to review my forgotten PR's on iceberg-rust (thanks btw), so I wanted to note that I have at least 3-4 PR's I want to land in this repo (datafusion-iceberg). I have avoided opening more PR's, to be considerate to the maintainers, but I'm wondering if that would be preferable as you are an approving reviewer here? Just to note, I also have at least 5 other PR's I want to open on iceberg-rust, as well. Let me know what works best for you. |
Thank you for the contributions!
I think the DataFusion community is pretty comfortable moving quickly on this repo, which is part of the reason the code moved out of Iceberg for the table provider. I would say keep firing off PRs, but anything you can do to signal ordering is helpful. PRs that aren't your highest priority could sit as draft until their dependencies merge, and EPIC issues track stacked PRs/larger tasks well.
That repo is a little trickier to manage, but we're trying over there, especially as it pertains to DataFusion-related PRs. Tag them with the Thanks again! |
new_with_predicate no longer takes a projection. When given, it had to repeat the schema's field names; when not, the schema had to be the full table schema, which was not checked, so a one-column schema with no projection built a scan that reported one column and returned two. The scan now always selects the names of its schema's fields, which reads the same columns as select_all for a full schema, so its batches match its schema by construction and the constructor cannot fail. This also fixes the catalog-backed provider, which pairs the schema it cached when built with the table as loaded on each scan: after a column was added and written to, a scan with no projection returned that column too. projection() keeps its signature and returns the schema's field names. The tests now compare full outputs rather than parts of them, check the pinned snapshot is read rather than the current one, and check each $snapshots row against the table's metadata.
The constructor and accessors went into a second inherent impl block; move them into the existing one, beside scan.
Which issue does this PR close?
What changes are included in this PR?
A distributed engine sends physical plans to other processes through a datafusion-proto
PhysicalExtensionCodec. To serialize a node, the codec has to name its type, read what it was built from, and build an equivalent node on the other side from those parts alone. For the Iceberg nodes, none of that is possible from outside this crate today:IcebergCommitExecandIcebergWriteExecarepub(crate), the modules of all four nodes are private, and none of them expose what they hold. This makes them public, together with read-only accessors for their parts, such that each public constructor takes only what the node's accessors return. It changes no behavior of existing plans.Newly exported from
datafusion_iceberg::physical_plan:IcebergCommitExec,IcebergWriteExecandIcebergMetadataScan, each with documented constructors and an accessor for what it holds:IcebergCommitExec::{new, table, catalog}.newno longer takes an Arrow schema: it was only shown in the verbose display and could not be read back from the node. The verbose display now shows the table's current schema, which is also the schema the commit reads the data files against.IcebergWriteExec::{new, table}IcebergMetadataScan::{new, provider}Providers:
IcebergMetadataTableProvider::{new, table, metadata_type}, and a re-export at the crate root. Its fields werepub(crate); they are now private behind the constructor.IcebergTableProvider::{try_new, catalog, table_ident}.try_newwaspub(crate); a codec needs it to rebuild the provider from its catalog and identifier.IcebergStaticTableProvider::{table, snapshot_id}.IcebergTableProvider::try_newoverlaps with #4; whichever lands second can drop it.IcebergTableScan::new_with_predicate. A codec can read a scan's pushed-downPredicate, but cannot turn it back into the DataFusion filters the existing constructor takes, so this constructor takes the predicate directly. It takes the scan's output schema, asschema()returns it, so a scan is rebuilt from its accessors alone, with no provider on the decoding side. The providers report pushdown asInexact, so DataFusion still applies the filters above the scan, and the predicate only lets Iceberg skip data files. The crate-privatenewnow projects DataFusion's column indices into the output schema and delegates to it; it returnsResult, turning an out-of-range index into an error where it used to panic (schema.project(..).unwrap()).IcebergTableScanreads its columns by name. The scan now always selects its schema's field names, where it usedselect_allfor a scan with no projection. For a full schema the two read the same columns: both resolve to the same field ids in iceberg-rust'scollect_scan_field_ids. Reading by name keeps the batches matching the schema by construction. It also fixes a mismatch inIcebergTableProvider, which pairs the schema it cached when built with the table as loaded on each scan: after a column was added and written to, a scan with no projection returned that column as well, beyond the provider's schema.projection()keeps its signature and now always returns the schema's field names, so EXPLAIN lists them for a full scan, where it used to printprojection:[], the same as an empty projection.IcebergCommitExecrefuses an input with more than one partition.executereads only partition 0 of its input and relies onrequired_input_distributionfor a coalesce the optimizer inserts. Built directly, as a codec does, over a multi-partition input, it would commit only the first partition's files and report success. With two input partitions of one file each, it committed 1 file and reported 100 rows instead of 200. It now returns an error before committing anything. Plans built byinsert_intoalways have a single-partition input, so they are unaffected.The constructor docs state what each node expects of its input (one partition for the commit; columns matched by name, and the partition values from
project_with_partitionfor partitioned tables, for the write).Are these changes tested?
test_plan_nodes_are_inspectable(integration):IcebergCatalogProvider, checks its identifier and catalog, and rebuilds it withtry_newfrom those parts.$snapshotsscan from its provider's parts, checks it returns the same rows, and checks each row against the table's metadata.expect!.test_scan_after_schema_evolution_reads_provider_columns(integration): after a column is added and written to, a scan through a provider built before the change returns only the provider's columns. It fails onselect_all.test_iceberg_commit_exec_rejects_multiple_input_partitions: two input partitions produce an error, and the table gets no snapshot.test_scan_rejects_out_of_range_projection: an out-of-range index is an error, not a panic.test_iceberg_commit_exec_empty_insert.new_with_predicateshows the rebuild round trip from the scan's accessors.cargo fmt --all -- --check,cargo clippy --workspace --locked --all-targets -- -D warnings,cargo test --workspace --lockedandRUSTDOCFLAGS="-D warnings" cargo doc --no-deps -p datafusion-icebergall pass locally.AI Disclosure