Skip to content

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

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

andygrove wants to merge 8 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 Spark 3.5 and later, Iceberg writes an unsorted partitioned table with 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. Taking the largest first frees the most memory per file closed. Of partitions holding as much, the one seen first closes first, so which files a write produces does not depend on HashMap order. Everything a fanout write reserves belongs to one of its partitions, so once they are all closed it reserves nothing, and the pool cannot fail it.
  • A closed partition's next rows open a new file with the writer properties its first file used, which the partition keeps.
  • The write then ends with more, smaller files. When every partition keeps getting rows, a write given 1/k of the memory it needs ends with roughly (k + 1) / 2 files per partition, however many partitions it writes to: about 4.3 at an eighth and 7.4 at a sixteenth.
  • 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, including how many files to expect, and its 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. A third checks the order partitions close in: 16 or 64 partitions taking turns through 32 batches, in an eighth of the memory the write needs, end with between four and five files per partition, and it asserts fewer than six. Closing every partition at once would leave seven, and closing them smallest first or in first-seen order seventeen to twenty, and the test fails under all three. 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 72 iceberg_write Rust tests and workspace clippy pass, and so do the Iceberg write action, detection, rewrite and system-function suites on Spark 4.1 (204 tests).

…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.
Comment thread native/core/src/execution/operators/iceberg_write.rs Outdated
Comment thread docs/source/contributor-guide/iceberg-writes.md Outdated
@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, ...)?

@sunchao sunchao left a comment

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.

Reviewed all seven changed files across all four commits against base f80042f78715e5a009352f1582f6c246a5f96131. The PR is not a draft. No introduced P1/P2 issues found within this review.

Read AGENTS.md, existing reviews, conversation comments and review threads. Applied review-comet-pr, review-comet-iceberg-write-pr and review-comet-memory-pr from .ai/skills/.

Summary

  • Prior state and problem: Native fanout writes retained a writer for every partition. When their buffers exceeded the task pool, reservation failure failed the task, and retrying could encounter the same limit.
  • Design approach: FanoutPartitions owns each partition’s feed, writer, memory counter and selected properties. After reservation refusal, it closes the largest partitions until the remaining reservation fits.
  • Correctness: Traced held-row flushing, repeated close/reopen, partition metadata, shared filename generation and cleanup ownership. PartitionFeed.finish() drains pending rows before closing, and reopened writers retain the selected properties. Compared Spark’s write/commit/abort lifecycle across supported source versions and iceberg-java’s fanout implementation. No introduced correctness issue met the P1/P2 bar.
  • Compatibility analysis: Checked Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0, plus Iceberg fanout contracts across 1.5.2–1.11.0. This change preserves planning gates, partition transforms, commit payloads and Java reconstruction of manifest metrics. It adds no version-dependent API calls or configuration defaults.
  • Key design decisions: One memory consumer remains responsible for the task. Per-partition counters include both open-file buffers and pending rows. Equal-size candidates use first-seen order, and final files remain sorted by path.
  • Implementation sketch: OpenFileMemory gains child counters, PartitionWriterBuilder.build_with() supports reopening with retained properties, and InnerWriter.reserve() performs close-and-retry. files_closed_early is exposed through the JVM SQL metrics. Rust tests, Spark tests and documentation cover the changed behavior.
  • Performance: Closing files releases writer buffers and lets constrained writes finish. The cost is additional files, writes and footer reads. Candidate sorting occurs only after reservation refusal and costs O(P log P) for P partitions. Existing per-batch pending-memory accounting remains linear. No measured throughput regression was established.
  • Design: The approach fits the required selective-close operation. The pinned iceberg-rust FanoutWriter exposes only whole-writer close and keeps its partition writers private, so a wrapper alone cannot implement this behavior.
  • Abstraction & complexity: FanoutPartition keeps related lifecycle state together, while the existing builder and feed implementations remain shared. The additional state has a clear role in accounting, reopening and collecting completed files. No concrete simplification meeting the reporting bar was identified.
  • Behavioral changes worth calling out: Fanout reservation refusal now produces more, smaller files instead of failing. Early close can end a file off the normal 1000-row grid and choose dictionary settings from fewer initial rows. These differences are documented. Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, pressure-driven closing is an intended change from keeping fanout writers open. Unpartitioned and clustered refusal behavior remains unchanged relative to this PR’s base.
  • Suggested improvements: No additional change met the P1/P2 reporting bar. The existing small-file discussion describes a real, documented tradeoff, but the evidence does not establish an unresolved P1/P2 defect. The other existing threads concern documentation and clarification.

Exact-head CI: run 37793476358 explicitly targets the reviewed SHA. It reports 58 successful checks and 11 skipped checks, with no failures. Logs confirm all 66 iceberg_write tests passed, including both early-close regressions and refusal cleanup. Spark 4.1 logs confirm the constrained-pool fanout row-parity test and unpartitioned failure/cleanup test passed. Iceberg 1.8, 1.9, 1.10 and 1.11 matrix jobs passed. Standalone Spark SQL matrices, macOS and benchmark checks were skipped.

Validation limits: The local focused Rust test command failed during dependency compilation because the shared filesystem ran out of space. No local tests executed. Its build artifacts were removed and the checkout remains clean. No local JVM suite, throughput benchmark or live S3/GCS write was run. Runtime conclusions rely on the verified exact-head CI results alongside source inspection.

Everything a fanout write reserves belongs to one of its partitions, so once
every partition holding memory is closed the resize is down to nothing, which
the pool cannot refuse. The contributor guide said a fanout write could still
fail with nothing left to close; only an unpartitioned or clustered write can.
FanoutPartition's key is used to open every file of the partition, not only to
write held rows out at close.
… holding the most

No test pinned the order in which a write the pool refuses closes its
partitions: closing every partition at once, or closing them smallest first or
in the order they were first seen, passed every test. The new test has 16 or 64
partitions take turns through 32 batches in an eighth of the memory the write
needs, and asserts fewer than six files per partition. Closing the largest first
writes between four and five, closing all at once seven, and the other two
orders seventeen.

The user guide now says how many files to expect, and that Iceberg's
rewrite_data_files compacts them.
@andygrove

andygrove commented Oct 9, 2026 •

Copy link
Copy Markdown
Member Author

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

Yes, that's the case where it happens, once the pool can't hold the write: every partition keeps getting rows, so a partition closed early opens a new file on its next rows. When the pool holds the write, nothing closes early and each partition gets one file, as before.

How many files per partition depends on how far short the pool falls, not on how many partitions there are. I measured it with rows taking turns across partitions as in your example (40 batches of 4096 rows with 100-byte payloads, about 18 MB), and the pool set to a fraction of the most the same write reserves with room to spare. Files per partition:

pool 16 partitions 64 partitions 256 partitions
1/2 1.69 1.66 1.66
1/4 2.62 2.55 2.56
1/8 4.38 4.31 4.27
1/16 7.44 7.36 7.29

So a pool 1/k of what the write needs gives roughly (k + 1) / 2 files per partition. Closing the partition holding the most is what keeps it that low, since that frees as much as one close can. At 1/8, closing every partition at once gives 7 files per partition, and closing them in first-seen order, which is smallest first here, gives 17 to 25. Nothing tested that order, so bb1788f adds a_fanout_write_short_of_memory_closes_the_partitions_holding_the_most, which fails under all three.

Before this PR the same write failed its task, and the retry failed the same way. iceberg-java's fanout writer keeps every file open until the task ends, so it writes one file per partition as long as the JVM heap holds all of them. Where the extra files matter, write.spark.fanout.enabled=false makes Spark sort each task's rows by partition, spilling as it needs to, and the task keeps one file open at a time. More off-heap memory, or rewrite_data_files afterwards, also work. The user guide now says roughly how many files to expect.

Avoiding the small files without more memory would take what Spark's own file writer does past spark.sql.maxConcurrentOutputFileWriters: sort the rest of the task's rows by partition and write them one partition at a time. That's a much bigger change, so I've filed it as #6818.

Closing partitions smallest first or in first-seen order writes 17 files per
partition with 16 partitions and 20 with 64; the comment only gave the first.

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

1. Drop the expect in FanoutPartition::write — let the type system prove the invariant

Avoiding unwrap()/expect() in favor of making the state unrepresentable: Option::insert sets the value and returns &mut T, so the "opened above" invariant becomes structural instead of a documented panic:

async fn write(
    &mut self,
    builder: &PartitionWriterBuilder,
    unit: RecordBatch,
) -> iceberg::Result<()> {
    let writer = match self.writer.as_mut() {
        Some(writer) => writer,
        None => {
            let properties = self
                .properties
                .get_or_insert_with(|| builder.properties.take(Some(self.key.data())))
                .clone();
            self.writer.insert(
                builder
                    .build_with(Some(self.key.clone()), properties, self.memory.clone())
                    .await?,
            )
        }
    };
    writer.write(unit).await
}

2. Guard OpenFileMemory::child() against a second level of nesting

The parent chain in OpenFileShare::set only propagates one level: a child() of a child would add to the middle counter but silently never reach the task-level total, under-counting the reservation. Today nobody does this, but nothing stops a future caller. Either a debug assertion:

/// A counter for one fanout partition's files, which also counts towards this one.
fn child(&self) -> OpenFileMemory {
    debug_assert!(
        self.parent.is_none(),
        "OpenFileMemory propagates one level; a child's child would not reach the task total"
    );
    OpenFileMemory {
        bytes: Arc::default(),
        parent: Some(Arc::clone(&self.bytes)),
    }
}

…or, more thoroughly, split the type: OpenFileMemory (task total, has child()) → PartitionMemory (has share() only), so the nesting rule is enforced at compile time.

3. Centralize the held_bytes diff bookkeeping

The held_before / subtract / re-add pattern is repeated at every call site and is easy to get wrong when a fourth site appears:

/// Reconciles the shared held-rows total after partition `i`'s feed changed,
/// given what it held before the change.
fn sync_held(&mut self, i: usize, held_before: usize) {
    self.held_bytes = self.held_bytes - held_before + self.partitions[i].feed.held_bytes();
}

used as:

let partition = &mut self.partitions[i];
let held_before = partition.feed.held_bytes();
let units = partition.feed.push(part, properties)?;
self.sync_held(i, held_before);

(close_early's -= partition.feed.held_bytes() is the same invariant from the other direction; a short doc comment on held_bytes stating "always kept equal to the sum of the feeds' held_bytes()" would tie the three sites together.)

4. Document why reserve goes through indexes rather than references

by_memory() returning Vec<usize> and reserve then re-indexing fanout.partitions[i] looks like an iterator-shaped hole; it's actually required (closing mutably borrows the partition while the ranking is alive, and the index doubles as the first-seen tie-break). One line on by_memory would save the next reader from trying to "fix" it:

/// The partitions holding memory, as indexes into `partitions`, the most first. Of two holding
/// as much, the one seen first comes first, so which files a write closes early does not
/// depend on `HashMap` order.
/// Indexes rather than references: `reserve` mutates (closes) each partition it picks while
/// the rest of the ranking must stay alive.

5. Minor: skip a copy in the test helper

payload_rows_from takes &[&str] only to .to_vec() it for StringArray::from, which accepts Vec<&str> directly — taking the Vec<&str> the caller already built avoids the copy and shortens the helper:

fn payload_rows_from(first: usize, regions: Vec<&str>, payload_bytes: usize) -> RecordBatch {
    let rows = regions.len();
    // …
    Arc::new(StringArray::from(regions)),
}

…ts counters in debug builds

- `FanoutPartition::write` takes the writer it opens from `Option::insert`,
  so it no longer needs an `expect`.
- `OpenFileMemory::child` asserts that it is called on the task's total,
  since a share counts only one level up.
- `FanoutPartitions::write` asserts that `held_bytes` is still the sum of
  what the feeds hold back, and the field's doc names the places that keep
  it so.
- `by_memory` says why it ranks partitions by index.
- `payload_rows_from` takes its regions by value instead of copying them.
@andygrove

Copy link
Copy Markdown
Member Author

Addressed in c29bb0a:

  1. Done: FanoutPartition::write takes the writer from Option::insert, so the expect is gone.

  2. Added the debug assertion, and child()'s doc now says only the task's total has children. I
    didn't split the type. The unpartitioned and clustered writers' files report straight to the
    task's total, so the total needs share() as much as a partition's counter does, and
    MeteredParquetWriterBuilder would have to take either type, for a rule with one caller.

  3. Only the push does the before-and-after diff. The release sets the total to zero once every feed
    has released, and close_early subtracts the partition's held rows before its close releases
    them. So sync_held would have one caller. It also takes &mut self while the loop after it
    still writes through partition, so that loop would have to look the partition up again.
    Instead, the doc comment on held_bytes now states the invariant and names the three places
    that keep it. A debug_assert_eq! at the end of write checks the total against the feeds,
    so a fourth place that misses it fails the tests: with close_early's subtraction removed,
    a_fanout_write_short_of_memory_closes_the_partitions_holding_the_most fails on it.

  4. Added to by_memory's doc.

  5. Done: payload_rows_from takes the Vec<&str>, as payload_batch_from did before this PR.

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

4 participants