Skip to content

feat: close fanout partitions early when the memory pool refuses the native Iceberg writer - #6773

Open
andygrove wants to merge 4 commits into
apache:mainfrom
andygrove:iceberg-fanout-close-largest
Open

andygrove wants to merge 4 commits into
apache:mainfrom
andygrove:iceberg-fanout-close-largest

Conversation

@andygrove

@andygrove andygrove commented Oct 8, 2026 •

Copy link
Copy Markdown
Member

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?

  • FanoutPartitions replaces iceberg-rust's FanoutWriter for 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, which FanoutWriter cannot.
  • Each fanout partition's files report what they hold to a child of the task's OpenFileMemory, so the writer can rank partitions by the memory they hold.
  • When the pool refuses the writer's reservation after a batch, InnerWriter::reserve 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 the partition keeps.
  • 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, with the same rows per partition as the write with room to spare, 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 its source, 219 files against 64 under the full pool, 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, rewrite and system-function suites pass on Spark 4.1 (204 tests), and so do the iceberg_write Rust tests and workspace clippy on Rust 1.99.

…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.
@andygrove

andygrove commented Oct 8, 2026 •

Copy link
Copy Markdown
Member Author

Why this copies iceberg-rust's FanoutWriter:

The fix needs two things that FanoutWriter can't do. iceberg-rust's FanoutWriter keeps its per-partition writers in a private HashMap and exposes only write(partition_key, batch) and close(self), which closes every partition's writer at once. When the pool refuses the reservation, this PR has to close one partition's writer while the task carries on, and has to know which partitions hold the most memory to choose them. Neither is possible from outside FanoutWriter.

So FanoutPartitions keeps the fanout writers itself: one writer per partition, opened on the partition's first rows and closed when the task ends, as in FanoutWriter. Each partition's record also holds its feed, the memory its open file holds and the properties its files use. That lets InnerWriter::reserve rank partitions by memory and close the largest, and lets a closed partition open its next file with the same properties.

iceberg-rust main has the same FanoutWriter API as the revision Comet pins. The closest upstream issue is apache/iceberg-rust#1744, which asks for a max_open_partitions cap. A cap on how many partitions are open doesn't bound memory, since what each partition holds depends on how many rows it gets. If upstream adds a way to close one partition's writer, the writers can go back to a FanoutWriter, and Comet can keep the rest of each partition's record on its own side.

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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this comment is stale?

`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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can the pool ever refuse write now?

@manuzhang

manuzhang commented Oct 8, 2026 •

Copy link
Copy Markdown
Member

When the pool refuses the writer's reservation after a batch, InnerWriter::reserve 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.

Could there be too many small files when the rows' partitions are interleaved (e.g. p0, p1, ... pN, p0, p1, ... pN, p0, p1, ...)?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Iceberg area:writer Native Parquet writer enhancement New feature or request run-iceberg-tests

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native fanout Iceberg writes fail when their open partitions outgrow the memory pool

2 participants