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]

Reply via email to