andygrove opened a new issue, #35: URL: https://github.com/apache/datafusion-iceberg/issues/35
### Describe the bug For partitioned tables, `insert_into` adds a `ProjectionExec` that computes the `_partition` column with `PartitionExpr`. `PartitionExpr` reads its source columns by position: `PartitionValueCalculator::calculate` indexes `batch.columns()` with positions precomputed from the table schema. But `PartitionExpr::children()` returns an empty list ([project.rs#L210-L212](https://github.com/apache/datafusion-iceberg/blob/a2bc9427d0659591b5f8122b90fec710fe2f5de6/crates/datafusion/src/physical_plan/project.rs#L210-L212)). Because the expression appears to reference no columns, DataFusion's `ProjectionPushdown` rule merges the partition projection into a `ProjectionExec` beneath it and leaves `PartitionExpr` unchanged. `PartitionExpr` then runs against the input batch of the lower projection, which has a different column layout. The data columns are still written correctly, but each data file gets partition values computed from another column. Reads that prune by partition then silently skip those files. ### To Reproduce ```sql -- pt: id INT NOT NULL, val INT NOT NULL, partitioned by identity(id) -- src (MemTable): a INT NOT NULL, b INT NOT NULL, rows (1, 100), (2, 200) INSERT INTO pt SELECT b AS id, a AS val FROM src; ``` Optimized physical plan: ``` IcebergCommitExec: table=ns.pt IcebergWriteExec: table=ns.pt ProjectionExec: expr=[b@1 as id, a@0 as val, iceberg_partition_values(id) as _partition] DataSourceExec: partitions=1, partition_sizes=[1] ``` The partition values come from column 0 of the source batch (`a`) instead of from `id` (`b`): - data files are written under `data/id=1` and `data/id=2` instead of `data/id=100` and `data/id=200` - `SELECT * FROM pt` returns `(100, 1)` and `(200, 2)` - `SELECT * FROM pt WHERE id = 100` returns no rows ### Expected behavior Partition values are computed from the columns being written, however the optimizer rearranges the projections. ### Additional context Any INSERT whose physical input ends in a `ProjectionExec` that renames, reorders or drops columns is exposed (column aliases, inserts from joins, ...). `VALUES` and `SELECT *` from a source with the table's column order work only because the positions happen to line up. The schema check in `project_with_partition` runs before physical optimization, so it doesn't catch this. Possible fixes: - expose the source columns as `Column` children, so projection rewrites update them - check the batch schema in `evaluate` and resolve source columns by name or field id - compute partition values inside `IcebergWriteExec` instead of in a `PhysicalExpr` -- 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]
