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.
Could there be too many small files when the rows' partitions are interleaved (e.g. p0, p1, ... pN, p0, p1, ... pN, p0, p1, ...)? |
sunchao
left a comment
There was a problem hiding this comment.
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:
FanoutPartitionsowns 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:
OpenFileMemorygains child counters,PartitionWriterBuilder.build_with()supports reopening with retained properties, andInnerWriter.reserve()performs close-and-retry.files_closed_earlyis 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)forPpartitions. 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
FanoutWriterexposes only whole-writer close and keeps its partition writers private, so a wrapper alone cannot implement this behavior. - Abstraction & complexity:
FanoutPartitionkeeps 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.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, 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.
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:
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 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, Avoiding the small files without more memory would take what Spark's own file writer does past |
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
left a comment
There was a problem hiding this comment.
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.
|
Addressed in c29bb0a:
|
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 Spark 3.5 and later, Iceberg writes an unsorted partitioned table with 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. 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 onHashMaporder. 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.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 pass, and so do the Iceberg write action, detection, rewrite and system-function suites on Spark 4.1 (204 tests).