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]

Reply via email to