andygrove opened a new pull request, #6773:
URL: https://github.com/apache/datafusion-comet/pull/6773

   ## Which issue does this PR close?
   
   Closes #6771.
   
   Part of #5644.
   
   ## Rationale for this change
   
   Since #6247, the native Iceberg writer's buffers count against the task's 
memory pool. A fanout write keeps a data file open for every partition a task 
writes to, and the writer cannot spill, so a task that writes to enough 
partitions fails with `Additional allocation failed for IcebergWriteExec` where 
iceberg-java's fanout writer, buffering on the JVM heap, succeeds. Spark's 
retry takes the same path, so the job fails. Review of #6664, which makes 
native writes the default, asked for this to be fixed first: on Iceberg 1.5+ an 
unsorted partitioned table uses the fanout writer by default.
   
   ## What changes are included in this PR?
   
   - `FanoutFiles` replaces iceberg-rust's `FanoutWriter` for fanout writes. It 
keeps a writer open per partition, as `FanoutWriter` does, and can also close 
one partition's writer before the task ends, which `FanoutWriter` cannot.
   - Each fanout file reports what it holds to its partition's `OpenFileMemory` 
as well as the task's.
   - When the pool refuses the writer's reservation after a batch, 
`run_write_task` writes out and closes the partitions holding the most memory, 
in their open file and in the rows their feed has not handed over, until the 
reservation fits. Of partitions holding as much, the one seen first closes 
first, so which files a write produces does not depend on `HashMap` order.
   - A closed partition's next rows open a new file with the writer properties 
its first file used, which `PartitionProperties::get` keeps for fanout 
partitions.
   - A `files_closed_early` metric on `CometIcebergWriteExec` counts the files 
closed early.
   - An unpartitioned or clustered write still fails when its one open file 
outgrows the pool.
   - Docs: the user guide's failure handling and accepted divergences, the 
contributor guide's Iceberg writes and memory management pages, and the Iceberg 
write review skill.
   
   ## How are these changes tested?
   
   - New Rust tests: a fanout write given half the pool it needs closes 
partitions and writes every row, into more files than partitions, and a fanout 
write whose memory is all in rows held back for the dictionary choice writes 
them out to fit. Both fail without the change, and the second also fails if 
held-back rows are left out of a partition's memory. The test of a write the 
pool cannot hold now uses an unpartitioned write, which has nothing to close.
   - `CometIcebergWriteActionSuite`: the 64-partition fanout write under a 4 
MiB pool now interleaves its partitions and succeeds, with the same rows as the 
write under the full pool, 219 files against 64, and a non-zero 
`files_closed_early`. A new test keeps the failure and cleanup coverage with an 
unpartitioned write whose file outgrows the pool.
   - The Iceberg write action, detection and rewrite suites pass on Spark 4.1 
(189 tests), and so do the `iceberg_write` Rust tests and workspace clippy on 
Rust 1.99.
   


-- 
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