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]
