Repository navigation
Conversation
…native Iceberg writer A native fanout write keeps a file open for every partition a task writes to, and its buffers count against the task's memory pool, so a task over many partitions could fail where iceberg-java's on-heap writer succeeded. When the pool refuses the writer's reservation, it now writes out and closes the partitions holding the most memory, in their open file or in rows not yet handed to it, until what is left fits. A closed partition's next rows open a new file with the same writer properties, so the write finishes with more, smaller files instead of failing. A files_closed_early metric counts them. Writes that keep one file open still fail when that file outgrows the pool.
…r memory pressure
|
Why this copies iceberg-rust's The fix needs two things that So iceberg-rust |
The control write is native too, so it could lose the same rows. Also fix a stale FanoutWriter mention in the user guide and a missing noun in the FanoutFiles doc comment.
… record FanoutPartitions holds every partition in first-seen order, each with its feed, its open writer, the OpenFileMemory child its files report to, and the properties its files are written with. That replaces FanoutFiles and FanoutFeeds and their join by key, the seen counter, PartitionProperties::get and the partition field threaded through the metered builder. InnerWriter::reserve now owns the close-until-it-fits loop, ranking partitions by index without cloning keys and tracking pending bytes without rescanning. WriteMetrics gets a constructor, the Rust memory tests share their setup, and the Scala test captures the write with captureWrite.
|
|
||
| /// One partition of a fanout write. | ||
| struct FanoutPartition { | ||
| /// Kept so held rows can still be written out at close. |
| `OpenFileMemory`, which is how `reserve` finds them. That is why the fanout path uses | ||
| `FanoutPartitions` rather than iceberg-rust's `FanoutWriter`, which cannot close one partition's | ||
| writer. A closed partition's next rows open a new file with the properties its first file used, | ||
| which the partition keeps. A write the pool still refuses, with nothing left to close, fails the |
There was a problem hiding this comment.
Can the pool ever refuse write now?
Could there be too many small files when the rows' partitions are interleaved (e.g. p0, p1, ... pN, p0, p1, ... pN, p0, p1, ...)? |
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 IcebergWriteExecwhere 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?
FanoutPartitionsreplaces iceberg-rust'sFanoutWriterfor fanout writes. It keeps every partition the task has seen, in first-seen order, each with its feed, its open writer and the writer properties its files use, and can close one partition's writer before the task ends, whichFanoutWritercannot.OpenFileMemory, so the writer can rank partitions by the memory they hold.InnerWriter::reservewrites 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 onHashMaporder.files_closed_earlymetric onCometIcebergWriteExeccounts the files closed early.How are these changes tested?
CometIcebergWriteActionSuite: the 64-partition fanout write under a 4 MiB pool now interleaves its partitions and succeeds, with the same rows as its source, 219 files against 64 under the full pool, and a non-zerofiles_closed_early. A new test keeps the failure and cleanup coverage with an unpartitioned write whose file outgrows the pool.iceberg_writeRust tests and workspace clippy on Rust 1.99.