andygrove opened a new issue, #38: URL: https://github.com/apache/datafusion-iceberg/issues/38
### Describe the bug With `write.datafusion.fanout.enabled=false`, `insert_into` places a `SortExec` on `_partition`, and a hash `RepartitionExec` on `_partition`, in front of `IcebergWriteExec` ([table/mod.rs#L200-L220](https://github.com/apache/datafusion-iceberg/blob/a2bc9427d0659591b5f8122b90fec710fe2f5de6/crates/datafusion/src/table/mod.rs#L200-L220)), because `ClusteredWriter` needs its input grouped by partition. `IcebergWriteExec` doesn't declare `required_input_ordering` or `required_input_distribution`. It reports `maintains_input_order() == [true]` but advertises no output ordering ([write.rs#L143-L151](https://github.com/apache/datafusion-iceberg/blob/a2bc9427d0659591b5f8122b90fec710fe2f5de6/crates/datafusion/src/physical_plan/write.rs#L143-L151)). DataFusion's `EnforceSorting` therefore treats the sort as unnecessary and removes it, and `EnforceDistribution` replaces the hash repartition. ### To Reproduce 1. Create `t6(id INT NOT NULL, cat INT NOT NULL)`, partitioned by `identity(cat)`, with `write.datafusion.fanout.enabled=false`. 2. Run `INSERT INTO t6 SELECT id, cat FROM src`, where `src` has 200 rows with `cat = id % 10` and `batch_size = 8`. Optimized plan with `target_partitions = 4`: ``` IcebergCommitExec: table=ns.t6 CoalescePartitionsExec IcebergWriteExec: table=ns.t6 ProjectionExec: expr=[id@0 as id, cat@1 as cat, iceberg_partition_values(cat) as _partition] RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1 DataSourceExec: partitions=1, partition_sizes=[1] ``` All 10 runs failed, at both `target_partitions = 1` and `4`: ``` The input is not sorted! Cannot write to partition that was previously closed: PartitionKey { ... } ``` ### Expected behavior The insert succeeds, with the writer receiving its input clustered by partition. ### Additional context - `test_insert_plan_fanout_disabled_has_sort` passes because it inspects the plan that `insert_into` returns, which hasn't been through physical optimization. - Declaring the requirements on `IcebergWriteExec` would let the optimizer keep or re-insert them: `required_input_ordering` on `_partition` when fanout is disabled, and a hash `required_input_distribution` on `_partition`. - With fanout enabled (the default), the hash repartition is also replaced by round-robin, so every write task writes every partition. The data is correct, but there are more and smaller files than intended. -- 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]
