From 53e8dca1dcd7214489e750d3ab07582c886d5200 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 07:12:19 -0600 Subject: [PATCH 1/4] feat: close fanout partitions early when the memory pool refuses the 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. --- .../contributor-guide/iceberg-writes.md | 20 +- .../user-guide/latest/iceberg-writes.md | 27 +- .../src/execution/operators/iceberg_write.rs | 501 +++++++++++++++--- .../sql/comet/CometIcebergWriteExec.scala | 9 +- .../comet/CometIcebergWriteActionSuite.scala | 84 ++- 5 files changed, 530 insertions(+), 111 deletions(-) diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index 1ef7539329a..0224a1494cb 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -411,16 +411,16 @@ the writer: run the write suites and the Iceberg Spark tests. The pin policy is ## Testing -| Suite | What it covers | -| ---------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `CometIcebergWriteActionSuite` | End-to-end writes through the split plan and the native writer: parity with iceberg-java, row-level DML, partition evolution, file order, cleanup on task and job failure, a fanout write outgrowing the memory pool, AQE re-planning. | -| `CometIcebergWriteDetectionSuite` | One case per eligibility rule, accepted and declined. | -| `CometIcebergSystemFunctionSuite` | Native `bucket`, `truncate`, `years`/`months`/`days`/`hours`, which keep a partitioned write's hash distribution and sort native end to end. | -| `CometIcebergRewriteActionSuite` | Iceberg's `rewrite_data_files` with the split plan and the native writer. | -| `IcebergWriteProtoTranslationSuite` | Translation of properties into `IcebergParquetWriteSettings` and the writer mode. | -| Rust tests in `iceberg_write.rs` and in `iceberg_partition_*.rs` | File rolling on the 1000-row grid, fanout order, clustered input checks, cleanup guard, manifest round trip, memory reservation, partition path rendering, partition values past `chrono`'s calendar. | -| Rust tests in `iceberg_dictionary.rs` | The per-column dictionary choice against answers recorded from parquet-mr, including columns either side of its cut-off. | -| `CometIcebergWriteBenchmark` | Native versus iceberg-java for unpartitioned, clustered, fanout and copy-on-write delete writes. It checks each arm's plan before timing it. | +| Suite | What it covers | +| ---------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | +| `CometIcebergWriteActionSuite` | End-to-end writes through the split plan and the native writer: parity with iceberg-java, row-level DML, partition evolution, file order, cleanup on task and job failure, writes outgrowing the memory pool, AQE re-planning. | +| `CometIcebergWriteDetectionSuite` | One case per eligibility rule, accepted and declined. | +| `CometIcebergSystemFunctionSuite` | Native `bucket`, `truncate`, `years`/`months`/`days`/`hours`, which keep a partitioned write's hash distribution and sort native end to end. | +| `CometIcebergRewriteActionSuite` | Iceberg's `rewrite_data_files` with the split plan and the native writer. | +| `IcebergWriteProtoTranslationSuite` | Translation of properties into `IcebergParquetWriteSettings` and the writer mode. | +| Rust tests in `iceberg_write.rs` and in `iceberg_partition_*.rs` | File rolling on the 1000-row grid, fanout order, clustered input checks, cleanup guard, manifest round trip, memory reservation, partition path rendering, partition values past `chrono`'s calendar. | +| Rust tests in `iceberg_dictionary.rs` | The per-column dictionary choice against answers recorded from parquet-mr, including columns either side of its cut-off. | +| `CometIcebergWriteBenchmark` | Native versus iceberg-java for unpartitioned, clustered, fanout and copy-on-write delete writes. It checks each arm's plan before timing it. | The Comet suites run against the Iceberg version each Spark profile pins in `spark/pom.xml`: 1.5.2 for Spark 3.4, 1.8.1 for 3.5, 1.10.0 for 4.0 and 4.2, and 1.11.0 for 4.1. Only the default profile diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index a04f6be61cf..5e09fb738ef 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -310,13 +310,21 @@ The native writer's buffers are charged to Comet's memory pool, the off-heap bud native operators draw on, where iceberg-java's buffers sit on the JVM heap. A fanout write keeps a data file open for every partition a task writes to. Each open file holds the row group it is writing in memory, up to `write.parquet.row-group-size-bytes`, and on S3 or GCS also the last row -group it flushed, which is uploaded once the next one is complete or the file closes. So a task -writing to many partitions needs memory in proportion to them. When the pool cannot grant it, the -task fails with a `CometNativeException` reading `Additional allocation failed for IcebergWriteExec` -instead of exceeding the executor's memory, and Spark retries it like any other task failure. Such a -write fits in less memory with the fanout writer disabled (`write.spark.fanout.enabled=false`): -Spark then sorts each task's rows by partition, and the task keeps one file open at a time. A -smaller row-group size also helps. Otherwise the write needs a larger `spark.memory.offHeap.size`. +group it flushed, which is uploaded once the next one is complete or the file closes. Each +partition also holds the rows that have not reached its file yet. So a task writing to many +partitions needs memory in proportion to them. When the pool cannot grant it, the task writes out +and closes the partitions holding the most memory until what is left fits, and a closed +partition's next rows open a new file. The write then finishes with more, smaller files than +iceberg-java's would, and the `files closed early to free memory` metric of its +`CometIcebergWrite` operator counts the files it closed early. Disabling the fanout writer +(`write.spark.fanout.enabled=false`) avoids them: Spark then sorts each task's rows by partition, +and the task keeps one file open at a time. A larger `spark.memory.offHeap.size` also helps. + +A write that keeps one file open, unpartitioned or clustered, has no partition to close. When that +file outgrows the pool, the task fails with a `CometNativeException` reading +`Additional allocation failed for IcebergWriteExec` instead of exceeding the executor's memory, and +Spark retries it like any other task failure. A smaller `write.parquet.row-group-size-bytes` or a +larger `spark.memory.offHeap.size` lets such a write fit. Partial results are never committed. The commit set is exactly the commit messages returned by successful tasks — a failed task contributes none — and if the job fails, the driver-side @@ -411,6 +419,11 @@ a data file but not what any reader computes from it: they may cross the target several grid steps apart, and the resulting files can differ in row count by an arbitrary number of 1000-row blocks. Do not rely on file-layout parity between the two writers; rely only on each file rolling on its own 1000-row boundary. +- A fanout write that the memory pool cannot hold closes partitions before the task ends (see + [Failure handling](#failure-handling)), where iceberg-java's fanout writer keeps every file open + until then. So a partition can have more, smaller files than iceberg-java writes, and a file + closed early ends off the 1000-row grid, as the last file of a task does. A partition closed + before its rows filled a first page makes its dictionary choice from the rows it has. - A fanout write lists a task's data files in file-path order, where iceberg-java lists them in its own `StructLikeMap` iteration order. Both are stable across runs, and neither is a documented ordering, but the manifest entry order becomes the scan-task order and so the row diff --git a/native/core/src/execution/operators/iceberg_write.rs b/native/core/src/execution/operators/iceberg_write.rs index ba9b08bd9da..4f12988cfbd 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -25,6 +25,8 @@ //! `ManifestWriter` against an in-memory `FileIO`. The JVM decodes the bytes with //! `ManifestFiles.read(...)` to recover the `DataFile`s for commit. +use std::cmp::Reverse; +use std::collections::hash_map::Entry; use std::collections::HashMap; use std::fmt; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -44,7 +46,7 @@ use datafusion::execution::TaskContext; use datafusion::physical_expr::{EquivalenceProperties, PhysicalExpr}; use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType}; use datafusion::physical_plan::metrics::{ - ExecutionPlanMetricsSet, MetricBuilder, MetricsSet, Time, + Count, ExecutionPlanMetricsSet, MetricBuilder, MetricsSet, Time, }; use datafusion::physical_plan::stream::RecordBatchStreamAdapter; use datafusion::physical_plan::{ @@ -69,9 +71,9 @@ use iceberg::writer::file_writer::{ FileWriter, FileWriterBuilder, ParquetWriter, ParquetWriterBuilder, }; use iceberg::writer::partitioning::clustered_writer::ClusteredWriter; -use iceberg::writer::partitioning::fanout_writer::FanoutWriter; use iceberg::writer::partitioning::unpartitioned_writer::UnpartitionedWriter; -use iceberg::writer::{CurrentFileStatus, IcebergWriterBuilder}; +use iceberg::writer::partitioning::PartitioningWriter; +use iceberg::writer::{CurrentFileStatus, IcebergWriter, IcebergWriterBuilder}; use iceberg::ErrorKind; #[cfg(test)] use parquet::arrow::PARQUET_FIELD_ID_META_KEY; @@ -123,17 +125,19 @@ struct PartitionWriterBuilder { storage: StorageWrites, } -#[async_trait::async_trait] -impl IcebergWriterBuilder for PartitionWriterBuilder { - type R = PartitionDataFileWriter; - - async fn build(&self, partition_key: Option) -> iceberg::Result { - let properties = self - .properties - .take(partition_key.as_ref().map(PartitionKey::data)); +impl PartitionWriterBuilder { + /// The data file writer for `partition_key`, writing with `properties`, whose files also + /// report what they hold to `partition`, if given. + async fn build_with( + &self, + partition_key: Option, + properties: WriterProperties, + partition: Option, + ) -> iceberg::Result { let parquet_builder = MeteredParquetWriterBuilder { inner: ParquetWriterBuilder::new(properties, Arc::clone(&self.schema)), open_files: self.open_files.clone(), + partition, storage: self.storage, }; let rolling_builder = RollingFileWriterBuilder::new( @@ -149,13 +153,28 @@ impl IcebergWriterBuilder for PartitionWriterBuilder { } } +/// How the unpartitioned and clustered writers open a partition's writer: once, with the +/// properties chosen for it, which are no longer needed after. +#[async_trait::async_trait] +impl IcebergWriterBuilder for PartitionWriterBuilder { + type R = PartitionDataFileWriter; + + async fn build(&self, partition_key: Option) -> iceberg::Result { + let properties = self + .properties + .take(partition_key.as_ref().map(PartitionKey::data)); + self.build_with(partition_key, properties, None).await + } +} + /// The Parquet writer properties each partition's files are written with: the table's, with the /// dictionary choice [`DictionaryChooser`] made from the partition's first rows. /// /// A partition's [`PartitionFeed`] records the choice before it hands the partition's first rows /// to the writer, and [`PartitionWriterBuilder`] takes it back out when that first write opens the -/// partition's writer. The partitioning writers open a partition's writer exactly once per task, -/// on its first write. +/// partition's writer. The unpartitioned and clustered writers open a partition's writer exactly +/// once per task, on its first write. A fanout partition whose file [`FanoutFiles`] closed early +/// opens another, so a fanout write reads the choice and leaves it in place. struct PartitionProperties { chooser: DictionaryChooser, /// Keyed by partition value; `None` is the whole task of an unpartitioned write. @@ -186,6 +205,21 @@ impl PartitionProperties { .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) .remove(&partition.cloned()); + self.or_base(chosen) + } + + /// The properties chosen for `partition`, kept for the next file the partition opens. + fn get(&self, partition: Option<&IcebergStruct>) -> WriterProperties { + let chosen = self + .chosen + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .get(&partition.cloned()) + .cloned(); + self.or_base(chosen) + } + + fn or_base(&self, chosen: Option) -> WriterProperties { debug_assert!( chosen.is_some(), "a partition's writer opened before its dictionary choice was made" @@ -194,12 +228,14 @@ impl PartitionProperties { } } -/// What a task's open data files hold in memory between them. +/// What a task's open data files hold in memory between them, or what one fanout partition's +/// open file holds. /// -/// iceberg-rust keeps each open file writer private inside its rolling and partitioning writers, -/// so the files report their own shares here through [`MeteredParquetWriter`], and `run_write_task` -/// reserves the total. A fanout write keeps one file open per partition, so this is what grows -/// with the partition count. +/// iceberg-rust keeps each open file writer private inside its rolling writer, so the files report +/// their own shares here through [`MeteredParquetWriter`], and `run_write_task` reserves the +/// task's total. A fanout write keeps one file open per partition, so the total grows with the +/// partition count, and [`FanoutFiles`] keeps each partition's too, to find the partitions to +/// close when the pool refuses the total. #[derive(Clone, Debug, Default)] struct OpenFileMemory(Arc); @@ -208,12 +244,23 @@ impl OpenFileMemory { self.0.load(Ordering::Relaxed) } - fn share(&self) -> OpenFileShare { + /// A newly opened file's share of this total and, for a fanout partition's file, of the + /// partition's. + fn share(&self, partition: Option) -> OpenFileShare { OpenFileShare { total: self.clone(), + partition, bytes: 0, } } + + fn change(&self, from: usize, to: usize) { + if to > from { + self.0.fetch_add(to - from, Ordering::Relaxed); + } else { + self.0.fetch_sub(from - to, Ordering::Relaxed); + } + } } /// One open file's part of [`OpenFileMemory`], given back when the file closes, or when it is @@ -221,19 +268,15 @@ impl OpenFileMemory { #[derive(Debug)] struct OpenFileShare { total: OpenFileMemory, + partition: Option, bytes: usize, } impl OpenFileShare { fn set(&mut self, bytes: usize) { - if bytes > self.bytes { - self.total - .0 - .fetch_add(bytes - self.bytes, Ordering::Relaxed); - } else { - self.total - .0 - .fetch_sub(self.bytes - bytes, Ordering::Relaxed); + self.total.change(self.bytes, bytes); + if let Some(partition) = &self.partition { + partition.change(self.bytes, bytes); } self.bytes = bytes; } @@ -246,11 +289,13 @@ impl Drop for OpenFileShare { } /// [`ParquetWriterBuilder`] whose files report what they hold in memory to the task's -/// [`OpenFileMemory`]. +/// [`OpenFileMemory`], and to their fanout partition's. #[derive(Clone, Debug)] struct MeteredParquetWriterBuilder { inner: ParquetWriterBuilder, open_files: OpenFileMemory, + /// The fanout partition the files belong to, if [`FanoutFiles`] tracks it. + partition: Option, /// How the files' storage takes the row groups they flush. storage: StorageWrites, } @@ -263,7 +308,7 @@ impl FileWriterBuilder for MeteredParquetWriterBuilder { let output_file = CountedOutput::wrap(output_file, self.storage, Arc::clone(&released)); Ok(MeteredParquetWriter { inner: self.inner.build(output_file).await?, - share: self.open_files.share(), + share: self.open_files.share(self.partition.clone()), released, }) } @@ -816,9 +861,11 @@ impl ExecutionPlan for IcebergWriteExec { partition: usize, context: Arc, ) -> DFResult { - // Time spent inside the iceberg-rust writer stack (write + close), excluding time spent - // waiting on the upstream input stream. Surfaced on the JVM exec's SQL metrics by name. - let write_time = MetricBuilder::new(&self.metrics).subset_time("write_time", partition); + let metrics = WriteMetrics { + write_time: MetricBuilder::new(&self.metrics).subset_time("write_time", partition), + files_closed_early: MetricBuilder::new(&self.metrics) + .counter("files_closed_early", partition), + }; // One consumer for the whole task, however many files it opens: a consumer per file // would shrink every other consumer's share of the fair pool as a fanout write widened. let reservation = MemoryConsumer::new(format!("IcebergWriteExec[{partition}]")) @@ -843,7 +890,7 @@ impl ExecutionPlan for IcebergWriteExec { writer_properties.as_ref().clone(), partition_id, task_attempt_id, - write_time, + metrics, reservation, ) .await?; @@ -909,14 +956,26 @@ impl DisplayAs for IcebergWriteExec { } } +/// What a write task reports, surfaced on the JVM exec's SQL metrics by name. +#[derive(Clone, Debug, Default)] +struct WriteMetrics { + /// Time spent inside the iceberg-rust writer stack (write + close), excluding time spent + /// waiting on the upstream input stream. + write_time: Time, + /// Data files a fanout write closed before the task ended, to give back the memory they held. + files_closed_early: Count, +} + /// One-shot per-task write coroutine. Builds the iceberg-rust writer stack, decorates each input /// batch with `PARQUET_FIELD_ID_META_KEY` metadata so iceberg-rust can match Arrow columns to -/// Iceberg field IDs, and routes through `UnpartitionedWriter`/`FanoutWriter`/`ClusteredWriter` +/// Iceberg field IDs, and routes through `UnpartitionedWriter`/[`FanoutFiles`]/`ClusteredWriter` /// depending on `writer_mode`. /// /// After every batch, `reservation` is resized to what the task's writers hold between batches: /// the open files' shares (see [`MeteredParquetWriter`]) and the rows waiting in their -/// [`PartitionFeed`]s. The writer cannot spill, so a reservation the pool refuses fails the task. +/// [`PartitionFeed`]s. When the pool refuses, a fanout write writes out and closes its +/// partitions, the one holding the most first, until what is left fits. The writer cannot spill, +/// so a reservation the pool still refuses with nothing left to close fails the task. /// /// On success the still-armed [`AbortOnDrop`] is returned along with the data files: the caller /// owns cleanup until the JVM acknowledges the output after recording its locations. @@ -930,7 +989,7 @@ async fn run_write_task( writer_properties: WriterProperties, partition_id: Option, task_attempt_id: Option, - write_time: Time, + metrics: WriteMetrics, reservation: MemoryReservation, ) -> DFResult<(Vec, AbortOnDrop)> { // The JVM exec wrapper stamps both ids per task; a missing id means the plan template was @@ -1012,7 +1071,7 @@ async fn run_write_task( PartitionFeed::new(None, slicer), ), (false, ProtoIcebergWriterMode::IcebergWriterFanout) => InnerWriter::Fanout( - FanoutWriter::new(data_file_builder), + FanoutFiles::new(data_file_builder), splitter()?, FanoutFeeds::default(), ), @@ -1033,12 +1092,30 @@ async fn run_write_task( if nests_floats { decorated = drop_unwritten_values(decorated)?; } - let timer = write_time.timer(); + let timer = metrics.write_time.timer(); writer.write(decorated, &properties).await?; + if let Err(mut refused) = + reservation.try_resize(open_files.bytes() + writer.pending_bytes()) + { + let mut fits = false; + for partition in writer.partitions_by_memory() { + writer.close_partition(&partition, &properties).await?; + metrics.files_closed_early.add(1); + match reservation.try_resize(open_files.bytes() + writer.pending_bytes()) { + Ok(()) => { + fits = true; + break; + } + Err(e) => refused = e, + } + } + if !fits { + return Err(refused); + } + } timer.done(); - reservation.try_resize(open_files.bytes() + writer.pending_bytes())?; } - let _timer = write_time.timer(); + let _timer = metrics.write_time.timer(); writer.close(&properties).await } .await; @@ -1057,7 +1134,80 @@ async fn run_write_task( } } -/// Enum-based dispatch over the three iceberg-rust partitioning writers, each paired with the +/// The fanout writer: like iceberg-rust's `FanoutWriter`, a data file writer open for every +/// partition the task has written to, except that a partition's file can be closed before the +/// task ends, to give back the memory it holds. The partition's next rows then open a new file, +/// with the same properties. +/// +/// iceberg-java's fanout writer keeps every file open until the task ends, its buffers growing on +/// the JVM heap. The native writer's count against the task's memory pool instead, so when the +/// pool refuses them, closing partitions early lets the write finish with more, smaller files +/// rather than fail (see [`InnerWriter::partitions_by_memory`]). +struct FanoutFiles { + builder: PartitionWriterBuilder, + /// Each open partition's writer, and what its open file holds. + open: HashMap, + /// The data files of the partitions closed early. + closed: Vec, +} + +impl FanoutFiles { + fn new(builder: PartitionWriterBuilder) -> Self { + Self { + builder, + open: HashMap::new(), + closed: Vec::new(), + } + } + + /// What `partition`'s open file holds, if it has one open. + fn held(&self, partition: &IcebergStruct) -> usize { + self.open + .get(partition) + .map_or(0, |(_, memory)| memory.bytes()) + } + + /// Closes `partition`'s file, if it has one open. + async fn close_partition(&mut self, partition: &IcebergStruct) -> iceberg::Result<()> { + if let Some((mut writer, _)) = self.open.remove(partition) { + self.closed.extend(writer.close().await?); + } + Ok(()) + } +} + +#[async_trait::async_trait] +impl PartitioningWriter for FanoutFiles { + async fn write( + &mut self, + partition_key: PartitionKey, + input: RecordBatch, + ) -> iceberg::Result<()> { + let (writer, _) = match self.open.entry(partition_key.data().clone()) { + Entry::Occupied(open) => open.into_mut(), + Entry::Vacant(vacant) => { + let memory = OpenFileMemory::default(); + let properties = self.builder.properties.get(Some(partition_key.data())); + let writer = self + .builder + .build_with(Some(partition_key), properties, Some(memory.clone())) + .await?; + vacant.insert((writer, memory)) + } + }; + writer.write(input).await + } + + async fn close(self) -> iceberg::Result> { + let mut data_files = self.closed; + for (_, (mut writer, _)) in self.open { + data_files.extend(writer.close().await?); + } + Ok(data_files) + } +} + +/// Enum-based dispatch over the three partitioning writers, each paired with the /// [`PartitionFeed`] state that holds a partition's first rows back for its dictionary choice and /// keeps its rolling writer on iceberg-java's row grid, and the two partitioned ones with the /// [`PartitionSplitter`] that cuts their input. Each variant takes the same builder so we can keep @@ -1066,11 +1216,7 @@ async fn run_write_task( enum InnerWriter { Unpartitioned(UnpartitionedWriter, PartitionFeed), /// The fanout writer keeps one file open per partition, so every partition is fed separately. - Fanout( - FanoutWriter, - PartitionSplitter, - FanoutFeeds, - ), + Fanout(FanoutFiles, PartitionSplitter, FanoutFeeds), /// The clustered writer closes a partition's file as soon as the next key arrives, so only the /// current key's held rows are live; they are written out before the switch. Clustered( @@ -1090,7 +1236,7 @@ impl InnerWriter { InnerWriter::Fanout(_, _, fanout) => fanout .feeds .values() - .map(|(_, feed)| feed.pending_bytes()) + .map(|(_, feed, _)| feed.pending_bytes()) .sum(), InnerWriter::Clustered(_, _, live) => { live.as_ref().map_or(0, |(_, feed)| feed.pending_bytes()) @@ -1098,6 +1244,58 @@ impl InnerWriter { } } + /// The partitions of a fanout write that hold memory, in their open file or in the rows their + /// feed has not handed over, 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. + /// + /// A task that the pool refuses closes them in this order until what is left fits. The + /// unpartitioned and clustered writers hold one file open, at most a row group, and have none + /// to close. + fn partitions_by_memory(&self) -> Vec { + let InnerWriter::Fanout(files, _, fanout) = self else { + return Vec::new(); + }; + let mut holding: Vec<_> = fanout + .feeds + .iter() + .map(|(partition, (_, feed, seen))| { + ( + files.held(partition) + feed.pending_bytes(), + *seen, + partition, + ) + }) + .filter(|(bytes, _, _)| *bytes > 0) + .collect(); + holding.sort_unstable_by_key(|(bytes, seen, _)| (Reverse(*bytes), *seen)); + holding + .into_iter() + .map(|(_, _, partition)| partition.clone()) + .collect() + } + + /// Writes out every row a fanout `partition` still holds and closes its file, so that it + /// holds nothing until its next rows open a new one. Closing at a partial unit ends the file + /// off iceberg-java's row grid, as the end of the task would. + async fn close_partition( + &mut self, + partition: &IcebergStruct, + properties: &PartitionProperties, + ) -> DFResult<()> { + let InnerWriter::Fanout(files, _, fanout) = self else { + return Ok(()); + }; + let (key, feed, _) = fanout + .feeds + .get_mut(partition) + .expect("a fanout partition holding memory has a feed"); + fanout.held_bytes -= feed.held_bytes(); + for unit in feed.finish(properties)? { + files.write(key.clone(), unit).await.map_err(iceberg_err)?; + } + files.close_partition(partition).await.map_err(iceberg_err) + } + /// Writes `batch` through the writer this task built, in the [`ROWS_DIVISOR`]-row units the /// rolling writer needs to re-check the target file size on iceberg-java's cadence. Rows are /// fed after partition splitting, so each partition's dictionary choice is made from its own @@ -1108,7 +1306,6 @@ impl InnerWriter { batch: RecordBatch, properties: &PartitionProperties, ) -> DFResult<()> { - use iceberg::writer::partitioning::PartitioningWriter; match self { InnerWriter::Unpartitioned(w, feed) => { for unit in feed.push(batch, properties)? { @@ -1118,10 +1315,12 @@ impl InnerWriter { } InnerWriter::Fanout(w, splitter, fanout) => { for (key, part) in splitter.split_groups(&batch)? { - let (key, feed) = fanout.feeds.entry(key.data().clone()).or_insert_with(|| { - let partition = Some(key.data().clone()); - (key, PartitionFeed::new(partition, splitter.slicer)) - }); + let seen = fanout.feeds.len(); + let (key, feed, _) = + fanout.feeds.entry(key.data().clone()).or_insert_with(|| { + let partition = Some(key.data().clone()); + (key, PartitionFeed::new(partition, splitter.slicer), seen) + }); let held_before = feed.held_bytes(); let units = feed.push(part, properties)?; fanout.held_bytes = fanout.held_bytes - held_before + feed.held_bytes(); @@ -1131,7 +1330,7 @@ impl InnerWriter { } // Every partition stays open, so the rows they hold back share one limit. if !properties.chooser.may_hold(fanout.held_bytes) { - for (key, feed) in fanout.feeds.values_mut() { + for (key, feed, _) in fanout.feeds.values_mut() { for unit in feed.release(properties)? { w.write(key.clone(), unit).await.map_err(iceberg_err)?; } @@ -1175,7 +1374,6 @@ impl InnerWriter { /// same at close: the leftovers land in the file that is open at that point, unless the /// pending target-size check rolls first -- exactly as they would have on the JVM path. async fn close(self, properties: &PartitionProperties) -> DFResult> { - use iceberg::writer::partitioning::PartitioningWriter; match self { InnerWriter::Unpartitioned(mut w, mut feed) => { for unit in feed.finish(properties)? { @@ -1184,13 +1382,13 @@ impl InnerWriter { w.close().await.map_err(iceberg_err) } InnerWriter::Fanout(mut w, _, fanout) => { - for (_, (key, mut feed)) in fanout.feeds { + for (_, (key, mut feed, _)) in fanout.feeds { for unit in feed.finish(properties)? { w.write(key.clone(), unit).await.map_err(iceberg_err)?; } } let mut data_files = w.close().await.map_err(iceberg_err)?; - // `FanoutWriter` holds its per-partition writers in a `HashMap` and `close` + // `FanoutFiles` holds its per-partition writers in a `HashMap` and `close` // iterates it directly, so the order it returns follows Rust's per-process // `RandomState` and differs on every run. That order becomes the manifest entry // order, then the scan-task order, then the row order of an unordered @@ -1737,8 +1935,9 @@ impl PartitionFeed { #[derive(Default)] struct FanoutFeeds { /// The `PartitionKey` is kept alongside each feed so held rows can still be written out at - /// close. - feeds: HashMap, + /// close, and so is how many partitions the task had seen before this one, which orders + /// partitions that hold as much memory. + feeds: HashMap, /// Bytes held back across all feeds. held_bytes: usize, } @@ -2832,7 +3031,7 @@ mod tests { writer_properties, Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await?; @@ -2932,7 +3131,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -2980,7 +3179,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -3047,8 +3246,8 @@ mod tests { /// A fanout task's data files must come back in a deterministic order. /// - /// iceberg-rust's `FanoutWriter` keeps its per-partition writers in a `HashMap` and - /// `close` iterates it directly, so the `DataFile` order it returns follows Rust's + /// `FanoutFiles` keeps its per-partition writers in a `HashMap` and `close` iterates it + /// directly, so the `DataFile` order it returns follows Rust's /// per-process `RandomState` -- a different order on every run. That order becomes the /// manifest entry order, which becomes the scan-task order, which becomes the row order /// of an unordered `SELECT *`. Iceberg's own @@ -4257,7 +4456,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4386,7 +4585,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4454,7 +4653,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4547,7 +4746,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4651,6 +4850,28 @@ mod tests { batches: Vec, writer_properties: WriterProperties, target_file_size_bytes: u64, + ) -> (TempDir, DFResult>) { + write_reporting_to( + pool, + schema, + writer_mode, + batches, + writer_properties, + target_file_size_bytes, + WriteMetrics::default(), + ) + .await + } + + /// [`write_reserving_from`], reporting to `metrics`. + async fn write_reporting_to( + pool: &Arc, + schema: Schema, + writer_mode: ProtoIcebergWriterMode, + batches: Vec, + writer_properties: WriterProperties, + target_file_size_bytes: u64, + metrics: WriteMetrics, ) -> (TempDir, DFResult>) { let temp_dir = TempDir::new().unwrap(); let spec = match writer_mode { @@ -4681,7 +4902,7 @@ mod tests { writer_properties, Some(0), Some(0), - Time::default(), + metrics, reservation, ) .await @@ -4894,16 +5115,93 @@ mod tests { assert_eq!(pool.reserved(), 0, "the write kept a reservation"); } - /// A write whose open files outgrow what the pool grants fails with the pool's error, - /// rather than holding memory nothing accounts for, and deletes the files it had opened. + /// Rows per partition in each data file of `data_files`, keyed by partition value. + fn rows_per_partition(data_files: &[DataFile]) -> HashMap { + let mut rows = HashMap::new(); + for file in data_files { + *rows.entry(file.partition().clone()).or_default() += file.record_count(); + } + rows + } + + /// A fanout write whose partitions outgrow what the pool grants writes out and closes the + /// ones holding the most until what is left fits, and carries on: a closed partition's + /// next rows open a new file. Every row still lands, in more files than partitions. #[tokio::test] - async fn a_write_the_pool_cannot_hold_fails_and_deletes_its_files() { + async fn a_fanout_write_the_pool_cannot_hold_closes_partitions_until_it_fits() { + // Four batches, each a unit for every one of 16 partitions, so a partition closed + // early gets more rows after it. + let batches = || { + (0..4) + .map(|batch| { + round_robin_batch_from(batch * 16 * ROWS_DIVISOR, 16 * ROWS_DIVISOR, 16) + }) + .collect::>() + }; + // One unit is one page, so every unit goes straight to its partition's file. + let properties = || { + WriterProperties::builder() + .set_data_page_row_count_limit(ROWS_DIVISOR) + .build() + }; + + let roomy = Arc::new(PeakMemoryPool::new(usize::MAX)); + let (_dir, written) = write_reserving_from( + &roomy, + iceberg_user_schema(), + ProtoIcebergWriterMode::IcebergWriterFanout, + batches(), + properties(), + TARGET_FILE_SIZE, + ) + .await; + assert_eq!( + record_counts(&written.unwrap()), + vec![4 * ROWS_DIVISOR as u64; 16] + ); + + let tight = Arc::new(PeakMemoryPool::new(roomy.peak() / 2)); + let metrics = WriteMetrics::default(); + let (dir, written) = write_reporting_to( + &tight, + iceberg_user_schema(), + ProtoIcebergWriterMode::IcebergWriterFanout, + batches(), + properties(), + TARGET_FILE_SIZE, + metrics.clone(), + ) + .await; + let data_files = written.expect("closing partitions early lets the write fit"); + assert!( + metrics.files_closed_early.value() > 0, + "no partition was closed early" + ); + assert!( + data_files.len() > 16, + "{} files for 16 partitions", + data_files.len() + ); + let rows = rows_per_partition(&data_files); + assert_eq!(rows.len(), 16); + assert!( + rows.values().all(|&rows| rows == 4 * ROWS_DIVISOR as u64), + "{rows:?}" + ); + assert_eq!(parquet_files_under(dir.path()).len(), data_files.len()); + assert_eq!(tight.reserved(), 0, "the write kept a reservation"); + } + + /// A partition's memory counts the rows its feed still holds back, and closing it writes + /// them out first. Here every partition's rows wait for their dictionary choice, so no + /// file is open when the pool refuses: the write only fits by writing partitions out. + #[tokio::test] + async fn closing_a_partition_early_writes_out_the_rows_it_holds_back() { let batches = || vec![round_robin_batch(16 * ROWS_DIVISOR, 16)]; let properties = || WriterProperties::builder().build(); - // With room to spare, every partition's unit reaches its own file. let roomy = Arc::new(PeakMemoryPool::new(usize::MAX)); - let (dir, written) = write_reserving_from( + let (_dir, written) = write_reserving_from( &roomy, iceberg_user_schema(), ProtoIcebergWriterMode::IcebergWriterFanout, @@ -4913,16 +5211,67 @@ mod tests { ) .await; assert_eq!(written.unwrap().len(), 16); - assert_eq!(parquet_files_under(dir.path()).len(), 16); let tight = Arc::new(PeakMemoryPool::new(roomy.peak() / 2)); - let (dir, written) = write_reserving_from( + let metrics = WriteMetrics::default(); + let (_dir, written) = write_reporting_to( &tight, iceberg_user_schema(), ProtoIcebergWriterMode::IcebergWriterFanout, batches(), properties(), TARGET_FILE_SIZE, + metrics.clone(), + ) + .await; + let data_files = written.expect("writing held rows out lets the write fit"); + assert!( + metrics.files_closed_early.value() > 0, + "no partition was closed early" + ); + assert_eq!(record_counts(&data_files), vec![ROWS_DIVISOR as u64; 16]); + assert_eq!(tight.reserved(), 0, "the write kept a reservation"); + } + + /// A write whose open file outgrows what the pool grants, with no partition it can close + /// early, fails with the pool's error rather than holding memory nothing accounts for, + /// and deletes the files it had opened. An unpartitioned write keeps its one file open. + #[tokio::test] + async fn a_write_the_pool_cannot_hold_fails_and_deletes_its_files() { + // One unit a batch, and one unit is one page, so the open file grows with every batch. + let batches = || { + (0..4) + .map(|batch| round_robin_batch_from(batch * ROWS_DIVISOR, ROWS_DIVISOR, 16)) + .collect::>() + }; + let properties = || { + WriterProperties::builder() + .set_data_page_row_count_limit(ROWS_DIVISOR) + .build() + }; + + // With room to spare, the rows reach one file. + let roomy = Arc::new(PeakMemoryPool::new(usize::MAX)); + let (dir, written) = write_reserving_from( + &roomy, + iceberg_user_schema(), + ProtoIcebergWriterMode::IcebergWriterUnpartitioned, + batches(), + properties(), + TARGET_FILE_SIZE, + ) + .await; + assert_eq!(written.unwrap().len(), 1); + assert_eq!(parquet_files_under(dir.path()).len(), 1); + + let tight = Arc::new(PeakMemoryPool::new(roomy.peak() / 2)); + let (dir, written) = write_reserving_from( + &tight, + iceberg_user_schema(), + ProtoIcebergWriterMode::IcebergWriterUnpartitioned, + batches(), + properties(), + TARGET_FILE_SIZE, ) .await; let error = written.expect_err("a write that outgrows the pool must fail"); @@ -4952,6 +5301,7 @@ mod tests { let builder = MeteredParquetWriterBuilder { inner: ParquetWriterBuilder::new(WriterProperties::builder().build(), schema), open_files: open_files.clone(), + partition: None, storage: StorageWrites::Through, }; let mut closed = builder @@ -4995,6 +5345,7 @@ mod tests { Arc::clone(&schema), ), open_files: open_files.clone(), + partition: None, storage, }; let location = format!("memory:/t/{storage:?}.parquet"); diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala index 83149ced77c..fedad6b3f92 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala @@ -126,15 +126,18 @@ case class CometIcebergWriteExec( // Names mirror Spark's stock `BatchWriteHelper` metrics (`numFiles` / `numOutputRows` / // `numOutputBytes`) so the Spark SQL UI shows the same row as it would for a non-Comet - // Iceberg write. `write_time` is pushed from the native operator by name (see - // `iceberg_write.rs`). + // Iceberg write. `write_time` and `files_closed_early` are pushed from the native operator by + // name (see `iceberg_write.rs`). override lazy val metrics: Map[String, SQLMetric] = Map( "numFiles" -> SQLMetrics.createMetric(sparkContext, "number of files written"), "numOutputRows" -> SQLMetrics.createMetric(sparkContext, "number of output rows"), "numOutputBytes" -> SQLMetrics.createSizeMetric(sparkContext, "written output"), "write_time" -> SQLMetrics.createNanoTimingMetric( sparkContext, - "time in native Iceberg writer")) + "time in native Iceberg writer"), + "files_closed_early" -> SQLMetrics.createMetric( + sparkContext, + "files closed early to free memory")) override def doExecute(): RDD[InternalRow] = { val columnarRdd = doExecuteColumnar() diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index 45a4fe597fd..e444e7fe0a2 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -3078,21 +3078,21 @@ class CometIcebergWriteActionSuite } } - // The native writer reserves what its open files hold from the task's memory pool. A fanout - // write keeps a file open for every partition, so over enough partitions it outgrows the pool - // and fails its task, which Spark can retry, instead of growing in memory that no budget - // accounts for until the executor is killed. - test("native acceleration: a fanout write that outgrows the memory pool fails its task") { + // The native writer reserves what its partitions hold from the task's memory pool. A fanout + // write keeps every partition it has seen open, so over enough partitions it outgrows the pool. + // It then writes out and closes the partitions holding the most, whose next rows open new + // files, and finishes with more, smaller files where it would otherwise fail its task. + test("native acceleration: a fanout write that outgrows the memory pool closes partitions") { assumeNativeAcceleration() withIcebergCatalog { _ => val session = spark import session.implicits._ - // One task writes 64 partitions of 1000 random 500-byte strings, a run of rows per - // partition. Each open file holds about 500 KB, so the pool of about 4 MiB (0.002 of the - // suite's 2 GiB off-heap size) is outgrown long before all 64 files are open. + // One task writes 64 partitions of 1000 random 500-byte strings, the partitions taking + // turns, so every partition keeps getting rows. That is about 32 MB, and the pool of about + // 4 MiB (0.002 of the suite's 2 GiB off-heap size) holds an eighth of it. val rows = 64000 (0 until rows) - .map(i => (i, s"r${i / 1000}", new scala.util.Random(i).alphanumeric.take(500).mkString)) + .map(i => (i, s"r${i % 64}", new scala.util.Random(i).alphanumeric.take(500).mkString)) .toDF("id", "region", "payload") .coalesce(1) .createOrReplaceTempView("fanout_oom_src") @@ -3109,10 +3109,63 @@ class CometIcebergWriteActionSuite } withNativeEnabled { + var plans = Seq.empty[SparkPlan] + // withSQLConf returns Unit before Spark 4.0, so the plans are kept from inside it. withSQLConf(CometConf.COMET_OFFHEAP_MEMORY_POOL_FRACTION.key -> "0.002") { - val (failedPlans, error) = captureFailedPlans(spark) { + plans = capturePlans(spark) { spark.sql(s"INSERT INTO $catalog.$ns.fanout_oom SELECT * FROM fanout_oom_src") } + } + val writers = + plans.flatMap(p => collectWithSubqueries(p) { case w: CometIcebergWriteExec => w }) + assert(writers.nonEmpty, s"the write did not run natively:\n${plans.mkString("\n--\n")}") + assert( + writers.map(_.metrics("files_closed_early").value).sum > 0, + "the write closed no partition early") + + // With the whole pool, every partition keeps its one file open. + spark.sql(s"INSERT INTO $catalog.$ns.fanout_oom_control SELECT * FROM fanout_oom_src") + assert(parquetFiles(dataDir("fanout_oom_control")).size == 64) + assert( + parquetFiles(dataDir("fanout_oom")).size > 64, + s"${parquetFiles(dataDir("fanout_oom")).size} files for 64 partitions") + def rowsOf(table: String) = + spark.sql(s"SELECT id, region, payload FROM $catalog.$ns.$table") + assert(rowsOf("fanout_oom").count() == rows) + assert(rowsOf("fanout_oom").exceptAll(rowsOf("fanout_oom_control")).isEmpty) + assert(rowsOf("fanout_oom_control").exceptAll(rowsOf("fanout_oom")).isEmpty) + } + } + } + + // A write that keeps one file open has no partition to close early, so when the file outgrows + // the pool the task fails, which Spark can retry, instead of growing in memory that no budget + // accounts for until the executor is killed. + test("native acceleration: a write whose open file outgrows the memory pool fails its task") { + assumeNativeAcceleration() + withIcebergCatalog { _ => + val session = spark + import session.implicits._ + // One task writes about 32 MB of random 500-byte strings into one file, whose row group + // outgrows the pool of about 4 MiB (0.002 of the suite's 2 GiB off-heap size). + val rows = 64000 + (0 until rows) + .map(i => (i, "r", new scala.util.Random(i).alphanumeric.take(500).mkString)) + .toDF("id", "region", "payload") + .coalesce(1) + .createOrReplaceTempView("file_oom_src") + Seq("file_oom", "file_oom_control").foreach { table => + spark.sql(s""" + CREATE TABLE $catalog.$ns.$table (id INT, region STRING, payload STRING) + USING iceberg + """) + } + + withNativeEnabled { + withSQLConf(CometConf.COMET_OFFHEAP_MEMORY_POOL_FRACTION.key -> "0.002") { + val (failedPlans, error) = captureFailedPlans(spark) { + spark.sql(s"INSERT INTO $catalog.$ns.file_oom SELECT * FROM file_oom_src") + } assert( error.toSeq .flatMap(exceptionChain) @@ -3125,23 +3178,22 @@ class CometIcebergWriteActionSuite collectWithSubqueries(p) { case w: CometIcebergWriteExec => w }.nonEmpty), s"the failed write did not run natively:\n${failedPlans.mkString("\n--\n")}") } - assert(countSnapshots("fanout_oom") == 0L, "the failed write must not commit") + assert(countSnapshots("file_oom") == 0L, "the failed write must not commit") assert( - parquetFiles(dataDir("fanout_oom")).isEmpty, - s"the failed task left data files behind: ${parquetFiles(dataDir("fanout_oom"))}") + parquetFiles(dataDir("file_oom")).isEmpty, + s"the failed task left data files behind: ${parquetFiles(dataDir("file_oom"))}") // The executor outlived the failure, and with the whole pool the same write succeeds. val controlPlans = capturePlans(spark) { - spark.sql(s"INSERT INTO $catalog.$ns.fanout_oom_control SELECT * FROM fanout_oom_src") + spark.sql(s"INSERT INTO $catalog.$ns.file_oom_control SELECT * FROM file_oom_src") } assert( controlPlans.exists(p => collectWithSubqueries(p) { case w: CometIcebergWriteExec => w }.nonEmpty), "the control write did not run natively") assert( - spark.sql(s"SELECT count(*) FROM $catalog.$ns.fanout_oom_control").head().getLong(0) + spark.sql(s"SELECT count(*) FROM $catalog.$ns.file_oom_control").head().getLong(0) == rows) - assert(parquetFiles(dataDir("fanout_oom_control")).size == 64) } } } From f517aeee8d8b6aecf35a3a24abdc1e87591df340 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 07:13:15 -0600 Subject: [PATCH 2/4] docs: describe how a native fanout write closes partitions early under memory pressure --- .../review-comet-iceberg-write-pr/SKILL.md | 4 +++- .../contributor-guide/iceberg-writes.md | 19 +++++++++++++------ .../contributor-guide/memory_management.md | 19 ++++++++++++------- 3 files changed, 28 insertions(+), 14 deletions(-) diff --git a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md index a835f45441a..0344f0fa267 100644 --- a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md +++ b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md @@ -156,7 +156,9 @@ Read the ownership table in the contributor guide before reviewing any change ne `MeteredParquetWriterBuilder`, and rows the writer holds outside iceberg-rust (the `PartitionFeed`s, or anything a new feed holds back) must be added to what `run_write_task` reserves. A builder that constructs `ParquetWriterBuilder` directly leaves its files out of the task's - memory reservation, so a wide fanout write grows past the pool instead of failing its task. + memory reservation, so a wide fanout write grows past the pool instead of closing + partitions early or failing its task. A fanout file must also report to its partition's + `OpenFileMemory`, or the partitions holding the most are not the ones closed. A storage scheme newly supported for writes needs its entry in `StorageWrites::for_location`: a file reports its flushed row groups until its storage has written them out, which a local file does at once and an object store only part by part. diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index 0224a1494cb..ad2dfbfbc04 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -223,7 +223,7 @@ per task, and `PartitionWriterBuilder` builds the rest of the stack once per par partition's Parquet writer properties: ```text -UnpartitionedWriter | FanoutWriter | ClusteredWriter +UnpartitionedWriter | FanoutFiles | ClusteredWriter -> PartitionWriterBuilder, per partition: ParquetWriterBuilder -> RollingFileWriterBuilder -> DataFileWriterBuilder ``` @@ -270,8 +270,9 @@ Points where Comet adapts iceberg-rust to match iceberg-java: - **Field ids and casting.** `decorate_batch_with_field_ids` casts each batch to the field-id-annotated Arrow schema derived from the Iceberg schema, with `safe: false`, so a type mismatch fails the task instead of writing NULLs. -- **Deterministic output order.** iceberg-rust's `FanoutWriter` closes its writers out of a - `HashMap`, so the fanout path sorts its `DataFile`s by path before returning them +- **Deterministic output order.** `FanoutFiles`, Comet's version of iceberg-rust's `FanoutWriter`, + closes its writers out of a `HashMap`, so the fanout path sorts its `DataFile`s by path before + returning them ([#5776](https://github.com/apache/datafusion-comet/issues/5776)). Manifest order becomes the row order of an unordered read, so any map iteration that reaches the output needs the same care. @@ -281,9 +282,15 @@ hold in memory. It hands `ParquetWriter` each file's `OutputFile` behind a `Coun writer counts the bytes that leave memory on their way to storage, and a file reports what it has written less those. When they leave depends on the storage (`StorageWrites`). After every batch `run_write_task` resizes the task's reservation to what the open files report plus the rows each -`PartitionFeed` holds, for the dictionary choice or for pacing, and a resize the pool refuses fails -the task. What that figure covers, and what it misses, is described under -[Native writers](memory_management.md#native-writers). +`PartitionFeed` holds, for the dictionary choice or for pacing. When the pool refuses a fanout +write, `run_write_task` writes out and closes the partitions holding the most, in their open file +and their feed, until the resize succeeds (`InnerWriter::partitions_by_memory`), and counts them in +the `files_closed_early` metric. That is why the fanout path uses `FanoutFiles` 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 `PartitionProperties::get` +keeps for it. A write the pool still refuses, with nothing left to close, fails the task, as an +unpartitioned or clustered write always does. What the reserved figure covers, and what it misses, +is described under [Native writers](memory_management.md#native-writers). `FileIO` comes from `load_file_io` in `iceberg_common.rs`, shared with the native scan. It picks the storage backend from the data location's scheme and wires in Comet's S3 credential bridge when diff --git a/docs/source/contributor-guide/memory_management.md b/docs/source/contributor-guide/memory_management.md index 7bd16d880f6..ae62167c9c8 100644 --- a/docs/source/contributor-guide/memory_management.md +++ b/docs/source/contributor-guide/memory_management.md @@ -378,9 +378,11 @@ An operator that never calls `try_grow` is invisible to the pool no matter how m ### Native writers Both native writers reserve what they hold between batches through a single consumer per task, -`ParquetWriterExec[N]` or `IcebergWriteExec[N]`, resized after every batch. Neither can spill, so -when the pool refuses a resize the task fails with a `CometNativeException` whose message starts -`Additional allocation failed for` and names the consumer. That is a task failure Spark can retry. +`ParquetWriterExec[N]` or `IcebergWriteExec[N]`, resized after every batch. Neither can spill. A +fanout Iceberg write can give memory back by closing partitions early, and does so when the pool +refuses a resize, as described below. Otherwise, when the pool refuses a resize, the task fails +with a `CometNativeException` whose message starts `Additional allocation failed for` and names the +consumer. That is a task failure Spark can retry. Unreserved, the same memory would count only toward the container limit, where exceeding it kills the executor. @@ -400,10 +402,13 @@ the executor. it has flushed since its last part was uploaded. At the default row-group size (`write.parquet.row-group-size-bytes`, 128 MiB) that is the last row group it flushed. The reservation also covers the rows each partition holds back, first for its dictionary choice and - then until they fill the 1000-row unit the rolling writer is fed in. It does not cover what - parquet-rs holds beyond the encoded size: dictionary hash tables, unencoded dictionary indices - and buffer capacity. Native writes decline Bloom filters today, and neither figure would include - them. + then until they fill the 1000-row unit the rolling writer is fed in. When the pool refuses a + fanout write's resize, the writer writes out and closes the partitions holding the most, file + and held rows together, until the resize succeeds, so the write ends with more files rather than + failing. Each file also reports to its partition's total, which is how the writer finds them. + The reservation does not cover what parquet-rs holds beyond the encoded size: dictionary hash + tables, unencoded dictionary indices and buffer capacity. Native writes decline Bloom filters + today, and neither figure would include them. The writers register one consumer per task rather than one per open file because every consumer registered with `fair_unified` lowers the share of every other consumer in the task. A consumer per From 68d4e340ad69010e6e157c196f4da5e3568eb7db Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 08:04:29 -0600 Subject: [PATCH 3/4] test: compare the early-closing fanout write with its source rows 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. --- docs/source/user-guide/latest/iceberg-writes.md | 5 ++--- native/core/src/execution/operators/iceberg_write.rs | 4 ++-- .../apache/comet/CometIcebergWriteActionSuite.scala | 11 ++++++----- 3 files changed, 10 insertions(+), 10 deletions(-) diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 5e09fb738ef..ce4bfd3dfa9 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -431,9 +431,8 @@ a data file but not what any reader computes from it: commit gives each file's rows, so the same rows can get different `_row_id` values from the two writers. The ids are unique either way, and across tasks iceberg-java's own assignment already depends on the order in which the tasks finish. Only the sorted order is reproducible on the - native path: iceberg-rust's `FanoutWriter` closes its per-partition writers out of a `HashMap`, - which under Rust's per-process `RandomState` would otherwise give a different order on every - run. Clustered and unpartitioned writes append in creation order on both paths and are + native path: the native fanout writer closes its per-partition writers out of a `HashMap`, which + under Rust's per-process `RandomState` would otherwise give a different order on every run. Clustered and unpartitioned writes append in creation order on both paths and are unaffected. - Compressed page bytes are implementation-defined: the codec and any explicit level are translated, but parquet-rs and parquet-mr embed different encoder implementations and diff --git a/native/core/src/execution/operators/iceberg_write.rs b/native/core/src/execution/operators/iceberg_write.rs index 4f12988cfbd..829660090e8 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -1140,8 +1140,8 @@ async fn run_write_task( /// with the same properties. /// /// iceberg-java's fanout writer keeps every file open until the task ends, its buffers growing on -/// the JVM heap. The native writer's count against the task's memory pool instead, so when the -/// pool refuses them, closing partitions early lets the write finish with more, smaller files +/// the JVM heap. The native writer's buffers count against the task's memory pool instead, so when +/// the pool refuses them, closing partitions early lets the write finish with more, smaller files /// rather than fail (see [`InnerWriter::partitions_by_memory`]). struct FanoutFiles { builder: PartitionWriterBuilder, diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index e444e7fe0a2..cf39fc47b3d 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -3129,11 +3129,12 @@ class CometIcebergWriteActionSuite assert( parquetFiles(dataDir("fanout_oom")).size > 64, s"${parquetFiles(dataDir("fanout_oom")).size} files for 64 partitions") - def rowsOf(table: String) = - spark.sql(s"SELECT id, region, payload FROM $catalog.$ns.$table") - assert(rowsOf("fanout_oom").count() == rows) - assert(rowsOf("fanout_oom").exceptAll(rowsOf("fanout_oom_control")).isEmpty) - assert(rowsOf("fanout_oom_control").exceptAll(rowsOf("fanout_oom")).isEmpty) + // Compared with the source rather than the control, which is native too and could lose + // the same rows. + val written = spark.sql(s"SELECT id, region, payload FROM $catalog.$ns.fanout_oom") + val source = spark.table("fanout_oom_src") + assert(written.count() == rows) + assert(written.exceptAll(source).isEmpty && source.exceptAll(written).isEmpty) } } } From 5d6288db67de09157f4ea75ca559631c7beeffdf Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 8 Oct 2026 08:33:23 -0600 Subject: [PATCH 4/4] refactor: keep each fanout partition's feed, writer and memory in one 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. --- .../contributor-guide/iceberg-writes.md | 25 +- .../user-guide/latest/iceberg-writes.md | 6 +- .../src/execution/operators/iceberg_write.rs | 671 ++++++++---------- .../comet/CometIcebergWriteActionSuite.scala | 23 +- 4 files changed, 341 insertions(+), 384 deletions(-) diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index ad2dfbfbc04..52f703b7e17 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -223,7 +223,7 @@ per task, and `PartitionWriterBuilder` builds the rest of the stack once per par partition's Parquet writer properties: ```text -UnpartitionedWriter | FanoutFiles | ClusteredWriter +UnpartitionedWriter | FanoutPartitions | ClusteredWriter -> PartitionWriterBuilder, per partition: ParquetWriterBuilder -> RollingFileWriterBuilder -> DataFileWriterBuilder ``` @@ -270,9 +270,9 @@ Points where Comet adapts iceberg-rust to match iceberg-java: - **Field ids and casting.** `decorate_batch_with_field_ids` casts each batch to the field-id-annotated Arrow schema derived from the Iceberg schema, with `safe: false`, so a type mismatch fails the task instead of writing NULLs. -- **Deterministic output order.** `FanoutFiles`, Comet's version of iceberg-rust's `FanoutWriter`, - closes its writers out of a `HashMap`, so the fanout path sorts its `DataFile`s by path before - returning them +- **Deterministic output order.** `FanoutPartitions`, Comet's version of iceberg-rust's + `FanoutWriter`, closes its files in an order that depends on how the rows arrived and on which + partitions it closed early, so the fanout path sorts its `DataFile`s by path before returning them ([#5776](https://github.com/apache/datafusion-comet/issues/5776)). Manifest order becomes the row order of an unordered read, so any map iteration that reaches the output needs the same care. @@ -282,14 +282,15 @@ hold in memory. It hands `ParquetWriter` each file's `OutputFile` behind a `Coun writer counts the bytes that leave memory on their way to storage, and a file reports what it has written less those. When they leave depends on the storage (`StorageWrites`). After every batch `run_write_task` resizes the task's reservation to what the open files report plus the rows each -`PartitionFeed` holds, for the dictionary choice or for pacing. When the pool refuses a fanout -write, `run_write_task` writes out and closes the partitions holding the most, in their open file -and their feed, until the resize succeeds (`InnerWriter::partitions_by_memory`), and counts them in -the `files_closed_early` metric. That is why the fanout path uses `FanoutFiles` 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 `PartitionProperties::get` -keeps for it. A write the pool still refuses, with nothing left to close, fails the task, as an -unpartitioned or clustered write always does. What the reserved figure covers, and what it misses, +`PartitionFeed` holds, for the dictionary choice or for pacing (`InnerWriter::reserve`). When the +pool refuses a fanout write, `reserve` writes out and closes the partitions holding the most, in +their open file and their feed, until the resize succeeds, and `run_write_task` counts them in the +`files_closed_early` metric. Each partition's files report to a child of the task's +`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 +task, as an unpartitioned or clustered write always does. What the reserved figure covers, and what it misses, is described under [Native writers](memory_management.md#native-writers). `FileIO` comes from `load_file_io` in `iceberg_common.rs`, shared with the native scan. It picks diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index ce4bfd3dfa9..60da2f74089 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -431,9 +431,9 @@ a data file but not what any reader computes from it: commit gives each file's rows, so the same rows can get different `_row_id` values from the two writers. The ids are unique either way, and across tasks iceberg-java's own assignment already depends on the order in which the tasks finish. Only the sorted order is reproducible on the - native path: the native fanout writer closes its per-partition writers out of a `HashMap`, which - under Rust's per-process `RandomState` would otherwise give a different order on every run. Clustered and unpartitioned writes append in creation order on both paths and are - unaffected. + native path: the order the native writer closes its files in depends on the order rows arrive in + and on which partitions it closes early to free memory. Clustered and unpartitioned writes append + in creation order on both paths and are unaffected. - Compressed page bytes are implementation-defined: the codec and any explicit level are translated, but parquet-rs and parquet-mr embed different encoder implementations and defaults (zstd default levels, LZ4 framing), so byte-identical output is not achievable even diff --git a/native/core/src/execution/operators/iceberg_write.rs b/native/core/src/execution/operators/iceberg_write.rs index 829660090e8..eb9b83d1f24 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -26,7 +26,6 @@ //! `ManifestFiles.read(...)` to recover the `DataFile`s for commit. use std::cmp::Reverse; -use std::collections::hash_map::Entry; use std::collections::HashMap; use std::fmt; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -126,18 +125,17 @@ struct PartitionWriterBuilder { } impl PartitionWriterBuilder { - /// The data file writer for `partition_key`, writing with `properties`, whose files also - /// report what they hold to `partition`, if given. + /// The data file writer for `partition_key`, writing with `properties`, whose files report + /// what they hold to `open_files`. async fn build_with( &self, partition_key: Option, properties: WriterProperties, - partition: Option, + open_files: OpenFileMemory, ) -> iceberg::Result { let parquet_builder = MeteredParquetWriterBuilder { inner: ParquetWriterBuilder::new(properties, Arc::clone(&self.schema)), - open_files: self.open_files.clone(), - partition, + open_files, storage: self.storage, }; let rolling_builder = RollingFileWriterBuilder::new( @@ -154,7 +152,7 @@ impl PartitionWriterBuilder { } /// How the unpartitioned and clustered writers open a partition's writer: once, with the -/// properties chosen for it, which are no longer needed after. +/// properties chosen for it. #[async_trait::async_trait] impl IcebergWriterBuilder for PartitionWriterBuilder { type R = PartitionDataFileWriter; @@ -163,7 +161,8 @@ impl IcebergWriterBuilder for PartitionWriterBuilder { let properties = self .properties .take(partition_key.as_ref().map(PartitionKey::data)); - self.build_with(partition_key, properties, None).await + self.build_with(partition_key, properties, self.open_files.clone()) + .await } } @@ -171,10 +170,9 @@ impl IcebergWriterBuilder for PartitionWriterBuilder { /// dictionary choice [`DictionaryChooser`] made from the partition's first rows. /// /// A partition's [`PartitionFeed`] records the choice before it hands the partition's first rows -/// to the writer, and [`PartitionWriterBuilder`] takes it back out when that first write opens the -/// partition's writer. The unpartitioned and clustered writers open a partition's writer exactly -/// once per task, on its first write. A fanout partition whose file [`FanoutFiles`] closed early -/// opens another, so a fanout write reads the choice and leaves it in place. +/// to the writer, and the partition takes it back out when that first write opens the partition's +/// writer, once per task. A fanout partition keeps what it took for the files it opens after +/// closing one early (see [`FanoutPartition`]). struct PartitionProperties { chooser: DictionaryChooser, /// Keyed by partition value; `None` is the whole task of an unpartitioned write. @@ -205,21 +203,6 @@ impl PartitionProperties { .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) .remove(&partition.cloned()); - self.or_base(chosen) - } - - /// The properties chosen for `partition`, kept for the next file the partition opens. - fn get(&self, partition: Option<&IcebergStruct>) -> WriterProperties { - let chosen = self - .chosen - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()) - .get(&partition.cloned()) - .cloned(); - self.or_base(chosen) - } - - fn or_base(&self, chosen: Option) -> WriterProperties { debug_assert!( chosen.is_some(), "a partition's writer opened before its dictionary choice was made" @@ -228,37 +211,37 @@ impl PartitionProperties { } } -/// What a task's open data files hold in memory between them, or what one fanout partition's -/// open file holds. +/// What a task's open data files hold in memory between them. /// /// iceberg-rust keeps each open file writer private inside its rolling writer, so the files report /// their own shares here through [`MeteredParquetWriter`], and `run_write_task` reserves the -/// task's total. A fanout write keeps one file open per partition, so the total grows with the -/// partition count, and [`FanoutFiles`] keeps each partition's too, to find the partitions to -/// close when the pool refuses the total. +/// total. A fanout write keeps one file open per partition, so this is what grows with the +/// partition count. Each fanout partition's files report to a [`child`](Self::child) of it, which +/// is how the write finds the partitions to close when the pool refuses the total. #[derive(Clone, Debug, Default)] -struct OpenFileMemory(Arc); +struct OpenFileMemory { + bytes: Arc, + /// The total this one is part of, for a fanout partition's. + parent: Option>, +} impl OpenFileMemory { fn bytes(&self) -> usize { - self.0.load(Ordering::Relaxed) + self.bytes.load(Ordering::Relaxed) } - /// A newly opened file's share of this total and, for a fanout partition's file, of the - /// partition's. - fn share(&self, partition: Option) -> OpenFileShare { - OpenFileShare { - total: self.clone(), - partition, - bytes: 0, + /// A counter for one fanout partition's files, which also counts towards this one. + fn child(&self) -> OpenFileMemory { + OpenFileMemory { + bytes: Arc::default(), + parent: Some(Arc::clone(&self.bytes)), } } - fn change(&self, from: usize, to: usize) { - if to > from { - self.0.fetch_add(to - from, Ordering::Relaxed); - } else { - self.0.fetch_sub(from - to, Ordering::Relaxed); + fn share(&self) -> OpenFileShare { + OpenFileShare { + total: self.clone(), + bytes: 0, } } } @@ -268,15 +251,17 @@ impl OpenFileMemory { #[derive(Debug)] struct OpenFileShare { total: OpenFileMemory, - partition: Option, bytes: usize, } impl OpenFileShare { fn set(&mut self, bytes: usize) { - self.total.change(self.bytes, bytes); - if let Some(partition) = &self.partition { - partition.change(self.bytes, bytes); + for total in std::iter::once(&self.total.bytes).chain(&self.total.parent) { + if bytes > self.bytes { + total.fetch_add(bytes - self.bytes, Ordering::Relaxed); + } else { + total.fetch_sub(self.bytes - bytes, Ordering::Relaxed); + } } self.bytes = bytes; } @@ -288,14 +273,12 @@ impl Drop for OpenFileShare { } } -/// [`ParquetWriterBuilder`] whose files report what they hold in memory to the task's -/// [`OpenFileMemory`], and to their fanout partition's. +/// [`ParquetWriterBuilder`] whose files report what they hold in memory to an +/// [`OpenFileMemory`]. #[derive(Clone, Debug)] struct MeteredParquetWriterBuilder { inner: ParquetWriterBuilder, open_files: OpenFileMemory, - /// The fanout partition the files belong to, if [`FanoutFiles`] tracks it. - partition: Option, /// How the files' storage takes the row groups they flush. storage: StorageWrites, } @@ -308,7 +291,7 @@ impl FileWriterBuilder for MeteredParquetWriterBuilder { let output_file = CountedOutput::wrap(output_file, self.storage, Arc::clone(&released)); Ok(MeteredParquetWriter { inner: self.inner.build(output_file).await?, - share: self.open_files.share(self.partition.clone()), + share: self.open_files.share(), released, }) } @@ -861,11 +844,7 @@ impl ExecutionPlan for IcebergWriteExec { partition: usize, context: Arc, ) -> DFResult { - let metrics = WriteMetrics { - write_time: MetricBuilder::new(&self.metrics).subset_time("write_time", partition), - files_closed_early: MetricBuilder::new(&self.metrics) - .counter("files_closed_early", partition), - }; + let metrics = WriteMetrics::new(&self.metrics, partition); // One consumer for the whole task, however many files it opens: a consumer per file // would shrink every other consumer's share of the fair pool as a fanout write widened. let reservation = MemoryConsumer::new(format!("IcebergWriteExec[{partition}]")) @@ -966,16 +945,25 @@ struct WriteMetrics { files_closed_early: Count, } +impl WriteMetrics { + fn new(metrics: &ExecutionPlanMetricsSet, partition: usize) -> Self { + Self { + write_time: MetricBuilder::new(metrics).subset_time("write_time", partition), + files_closed_early: MetricBuilder::new(metrics) + .counter("files_closed_early", partition), + } + } +} + /// One-shot per-task write coroutine. Builds the iceberg-rust writer stack, decorates each input /// batch with `PARQUET_FIELD_ID_META_KEY` metadata so iceberg-rust can match Arrow columns to -/// Iceberg field IDs, and routes through `UnpartitionedWriter`/[`FanoutFiles`]/`ClusteredWriter` -/// depending on `writer_mode`. +/// Iceberg field IDs, and routes through `UnpartitionedWriter`/[`FanoutPartitions`]/ +/// `ClusteredWriter` depending on `writer_mode`. /// -/// After every batch, `reservation` is resized to what the task's writers hold between batches: -/// the open files' shares (see [`MeteredParquetWriter`]) and the rows waiting in their -/// [`PartitionFeed`]s. When the pool refuses, a fanout write writes out and closes its -/// partitions, the one holding the most first, until what is left fits. The writer cannot spill, -/// so a reservation the pool still refuses with nothing left to close fails the task. +/// After every batch, `reservation` is resized to what the task's writers hold between batches +/// (see [`InnerWriter::reserve`]), and a fanout write the pool refuses closes partitions early +/// until what is left fits. The writer cannot spill, so a reservation the pool still refuses +/// fails the task. /// /// On success the still-armed [`AbortOnDrop`] is returned along with the data files: the caller /// owns cleanup until the JVM acknowledges the output after recording its locations. @@ -1070,11 +1058,9 @@ async fn run_write_task( UnpartitionedWriter::new(data_file_builder), PartitionFeed::new(None, slicer), ), - (false, ProtoIcebergWriterMode::IcebergWriterFanout) => InnerWriter::Fanout( - FanoutFiles::new(data_file_builder), - splitter()?, - FanoutFeeds::default(), - ), + (false, ProtoIcebergWriterMode::IcebergWriterFanout) => { + InnerWriter::Fanout(FanoutPartitions::new(data_file_builder), splitter()?) + } (false, ProtoIcebergWriterMode::IcebergWriterClustered) => { InnerWriter::Clustered(ClusteredWriter::new(data_file_builder), splitter()?, None) } @@ -1094,25 +1080,10 @@ async fn run_write_task( } let timer = metrics.write_time.timer(); writer.write(decorated, &properties).await?; - if let Err(mut refused) = - reservation.try_resize(open_files.bytes() + writer.pending_bytes()) - { - let mut fits = false; - for partition in writer.partitions_by_memory() { - writer.close_partition(&partition, &properties).await?; - metrics.files_closed_early.add(1); - match reservation.try_resize(open_files.bytes() + writer.pending_bytes()) { - Ok(()) => { - fits = true; - break; - } - Err(e) => refused = e, - } - } - if !fits { - return Err(refused); - } - } + let closed = writer + .reserve(&reservation, &open_files, &properties) + .await?; + metrics.files_closed_early.add(closed); timer.done(); } let _timer = metrics.write_time.timer(); @@ -1134,75 +1105,199 @@ async fn run_write_task( } } -/// The fanout writer: like iceberg-rust's `FanoutWriter`, a data file writer open for every -/// partition the task has written to, except that a partition's file can be closed before the -/// task ends, to give back the memory it holds. The partition's next rows then open a new file, -/// with the same properties. +/// A fanout write: like iceberg-rust's `FanoutWriter`, a data file writer open for every +/// partition the task has written to, except that a partition can be closed before the task ends, +/// to give back the memory it holds. Its next rows then open a new file, with the same properties. +/// `FanoutWriter` keeps its writers private and closes them only all at once, so the fanout path +/// keeps its own. /// /// iceberg-java's fanout writer keeps every file open until the task ends, its buffers growing on /// the JVM heap. The native writer's buffers count against the task's memory pool instead, so when /// the pool refuses them, closing partitions early lets the write finish with more, smaller files -/// rather than fail (see [`InnerWriter::partitions_by_memory`]). -struct FanoutFiles { +/// rather than fail (see [`InnerWriter::reserve`]). +struct FanoutPartitions { builder: PartitionWriterBuilder, - /// Each open partition's writer, and what its open file holds. - open: HashMap, + /// Every partition the task has seen, in the order it first saw them. + partitions: Vec, + /// Where each partition value sits in `partitions`. + index: HashMap, + /// Bytes the feeds hold back for their dictionary choice, across all partitions. Every + /// partition stays open, so the rows they hold back share one limit: once they reach it, every + /// partition still holding makes its choice from the rows it has. A task fanning out to many + /// partitions that each get less than a page would otherwise hold all of its rows back, + /// uncompressed, until it closed. + held_bytes: usize, /// The data files of the partitions closed early. closed: Vec, } -impl FanoutFiles { +/// One partition of a fanout write. +struct FanoutPartition { + /// Kept so held rows can still be written out at close. + key: PartitionKey, + feed: PartitionFeed, + /// The writer of the partition's open file, if it has one open. + writer: Option, + /// What the open file holds, as part of the task's [`OpenFileMemory`]. + memory: OpenFileMemory, + /// The properties the partition's files are written with, taken when its first file opens and + /// kept for the files it opens after closing one early. + properties: Option, +} + +impl FanoutPartition { + /// What the partition holds: its open file and the rows its feed has not handed over. + fn bytes(&self) -> usize { + self.memory.bytes() + self.feed.pending_bytes() + } + + /// Writes `unit` to the partition's open file, opening one first if it has none. + async fn write( + &mut self, + builder: &PartitionWriterBuilder, + unit: RecordBatch, + ) -> iceberg::Result<()> { + if self.writer.is_none() { + let properties = self + .properties + .get_or_insert_with(|| builder.properties.take(Some(self.key.data()))) + .clone(); + self.writer = Some( + builder + .build_with(Some(self.key.clone()), properties, self.memory.clone()) + .await?, + ); + } + self.writer + .as_mut() + .expect("the partition's file was opened above") + .write(unit) + .await + } + + /// Writes out every row the partition still holds back and closes its file, returning that + /// file's data files. A file closed at a partial unit ends off iceberg-java's row grid, as the + /// last file of a task does. + async fn close( + &mut self, + builder: &PartitionWriterBuilder, + properties: &PartitionProperties, + ) -> DFResult> { + for unit in self.feed.finish(properties)? { + self.write(builder, unit).await.map_err(iceberg_err)?; + } + match self.writer.take() { + Some(mut writer) => writer.close().await.map_err(iceberg_err), + None => Ok(Vec::new()), + } + } +} + +impl FanoutPartitions { fn new(builder: PartitionWriterBuilder) -> Self { Self { builder, - open: HashMap::new(), + partitions: Vec::new(), + index: HashMap::new(), + held_bytes: 0, closed: Vec::new(), } } - /// What `partition`'s open file holds, if it has one open. - fn held(&self, partition: &IcebergStruct) -> usize { - self.open - .get(partition) - .map_or(0, |(_, memory)| memory.bytes()) + /// Memory held by the rows waiting in the feeds, for every partition the task has seen. + fn pending_bytes(&self) -> usize { + self.partitions + .iter() + .map(|partition| partition.feed.pending_bytes()) + .sum() } - /// Closes `partition`'s file, if it has one open. - async fn close_partition(&mut self, partition: &IcebergStruct) -> iceberg::Result<()> { - if let Some((mut writer, _)) = self.open.remove(partition) { - self.closed.extend(writer.close().await?); + /// Writes `batch`, feeding each of its partitions separately. + async fn write( + &mut self, + splitter: &PartitionSplitter, + batch: &RecordBatch, + properties: &PartitionProperties, + ) -> DFResult<()> { + for (key, part) in splitter.split_groups(batch)? { + let next = self.partitions.len(); + let i = *self.index.entry(key.data().clone()).or_insert(next); + if i == next { + self.partitions.push(FanoutPartition { + feed: PartitionFeed::new(Some(key.data().clone()), splitter.slicer), + key, + writer: None, + memory: self.builder.open_files.child(), + properties: None, + }); + } + let partition = &mut self.partitions[i]; + let held_before = partition.feed.held_bytes(); + let units = partition.feed.push(part, properties)?; + self.held_bytes = self.held_bytes - held_before + partition.feed.held_bytes(); + for unit in units { + partition + .write(&self.builder, unit) + .await + .map_err(iceberg_err)?; + } + } + if !properties.chooser.may_hold(self.held_bytes) { + for partition in &mut self.partitions { + for unit in partition.feed.release(properties)? { + partition + .write(&self.builder, unit) + .await + .map_err(iceberg_err)?; + } + } + self.held_bytes = 0; } Ok(()) } -} -#[async_trait::async_trait] -impl PartitioningWriter for FanoutFiles { - async fn write( - &mut self, - partition_key: PartitionKey, - input: RecordBatch, - ) -> iceberg::Result<()> { - let (writer, _) = match self.open.entry(partition_key.data().clone()) { - Entry::Occupied(open) => open.into_mut(), - Entry::Vacant(vacant) => { - let memory = OpenFileMemory::default(); - let properties = self.builder.properties.get(Some(partition_key.data())); - let writer = self - .builder - .build_with(Some(partition_key), properties, Some(memory.clone())) - .await?; - vacant.insert((writer, memory)) - } - }; - writer.write(input).await + /// 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. + fn by_memory(&self) -> Vec { + let mut holding: Vec<(usize, usize)> = self + .partitions + .iter() + .enumerate() + .map(|(i, partition)| (partition.bytes(), i)) + .filter(|&(bytes, _)| bytes > 0) + .collect(); + holding.sort_unstable_by_key(|&(bytes, i)| (Reverse(bytes), i)); + holding.into_iter().map(|(_, i)| i).collect() } - async fn close(self) -> iceberg::Result> { - let mut data_files = self.closed; - for (_, (mut writer, _)) in self.open { - data_files.extend(writer.close().await?); - } + /// Closes the partition at `i` before the task ends, keeping its data files for the output. + async fn close_early(&mut self, i: usize, properties: &PartitionProperties) -> DFResult<()> { + let partition = &mut self.partitions[i]; + self.held_bytes -= partition.feed.held_bytes(); + let data_files = partition.close(&self.builder, properties).await?; + self.closed.extend(data_files); + Ok(()) + } + + /// Closes every partition and returns the task's data files in file-path order. + async fn close(mut self, properties: &PartitionProperties) -> DFResult> { + let mut data_files = std::mem::take(&mut self.closed); + for partition in &mut self.partitions { + data_files.extend(partition.close(&self.builder, properties).await?); + } + // The files come back in the order the partitions closed, which depends on the order the + // rows arrived in and on which partitions closed early. That order becomes the manifest + // entry order, then the scan-task order, then the row order of an unordered `SELECT *`, + // which Iceberg's own + // `TestMetadataTablesWithPartitionEvolution.testPartitionColumnNamedPartition` compares + // positionally against what iceberg-java wrote (apache/datafusion-comet#5776). + // + // Sorting by path is what makes the task's output reproducible. Path order sorts by + // partition directory, then by the file name counter within a partition, and needs nothing + // from the write order. The clustered and unpartitioned writers append in creation order + // and are already deterministic, so only the fanout path sorts. + data_files.sort_unstable_by(|a, b| a.file_path().cmp(b.file_path())); Ok(data_files) } } @@ -1216,7 +1311,7 @@ impl PartitioningWriter for FanoutFiles { enum InnerWriter { Unpartitioned(UnpartitionedWriter, PartitionFeed), /// The fanout writer keeps one file open per partition, so every partition is fed separately. - Fanout(FanoutFiles, PartitionSplitter, FanoutFeeds), + Fanout(FanoutPartitions, PartitionSplitter), /// The clustered writer closes a partition's file as soon as the next key arrives, so only the /// current key's held rows are live; they are written out before the switch. Clustered( @@ -1233,67 +1328,42 @@ impl InnerWriter { fn pending_bytes(&self) -> usize { match self { InnerWriter::Unpartitioned(_, feed) => feed.pending_bytes(), - InnerWriter::Fanout(_, _, fanout) => fanout - .feeds - .values() - .map(|(_, feed, _)| feed.pending_bytes()) - .sum(), + InnerWriter::Fanout(fanout, _) => fanout.pending_bytes(), InnerWriter::Clustered(_, _, live) => { live.as_ref().map_or(0, |(_, feed)| feed.pending_bytes()) } } } - /// The partitions of a fanout write that hold memory, in their open file or in the rows their - /// feed has not handed over, 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. - /// - /// A task that the pool refuses closes them in this order until what is left fits. The - /// unpartitioned and clustered writers hold one file open, at most a row group, and have none - /// to close. - fn partitions_by_memory(&self) -> Vec { - let InnerWriter::Fanout(files, _, fanout) = self else { - return Vec::new(); - }; - let mut holding: Vec<_> = fanout - .feeds - .iter() - .map(|(partition, (_, feed, seen))| { - ( - files.held(partition) + feed.pending_bytes(), - *seen, - partition, - ) - }) - .filter(|(bytes, _, _)| *bytes > 0) - .collect(); - holding.sort_unstable_by_key(|(bytes, seen, _)| (Reverse(*bytes), *seen)); - holding - .into_iter() - .map(|(_, _, partition)| partition.clone()) - .collect() - } - - /// Writes out every row a fanout `partition` still holds and closes its file, so that it - /// holds nothing until its next rows open a new one. Closing at a partial unit ends the file - /// off iceberg-java's row grid, as the end of the task would. - async fn close_partition( + /// Resizes `reservation` to what the task's writers hold between batches: the open files' + /// shares (`open_files`, see [`MeteredParquetWriter`]) and the rows waiting in their + /// [`PartitionFeed`]s. When the pool refuses, a fanout write writes out and closes its + /// partitions, the one holding the most first, until what is left fits, and this returns how + /// many it closed. The unpartitioned and clustered writers hold one file open, at most a row + /// group, and cannot close it early, so the pool's refusal stands. + async fn reserve( &mut self, - partition: &IcebergStruct, + reservation: &MemoryReservation, + open_files: &OpenFileMemory, properties: &PartitionProperties, - ) -> DFResult<()> { - let InnerWriter::Fanout(files, _, fanout) = self else { - return Ok(()); + ) -> DFResult { + let mut pending = self.pending_bytes(); + let mut refused = match reservation.try_resize(open_files.bytes() + pending) { + Ok(()) => return Ok(0), + Err(refused) => refused, }; - let (key, feed, _) = fanout - .feeds - .get_mut(partition) - .expect("a fanout partition holding memory has a feed"); - fanout.held_bytes -= feed.held_bytes(); - for unit in feed.finish(properties)? { - files.write(key.clone(), unit).await.map_err(iceberg_err)?; + let InnerWriter::Fanout(fanout, _) = self else { + return Err(refused); + }; + for (closed, i) in fanout.by_memory().into_iter().enumerate() { + pending -= fanout.partitions[i].feed.pending_bytes(); + fanout.close_early(i, properties).await?; + match reservation.try_resize(open_files.bytes() + pending) { + Ok(()) => return Ok(closed + 1), + Err(e) => refused = e, + } } - files.close_partition(partition).await.map_err(iceberg_err) + Err(refused) } /// Writes `batch` through the writer this task built, in the [`ROWS_DIVISOR`]-row units the @@ -1313,31 +1383,8 @@ impl InnerWriter { } Ok(()) } - InnerWriter::Fanout(w, splitter, fanout) => { - for (key, part) in splitter.split_groups(&batch)? { - let seen = fanout.feeds.len(); - let (key, feed, _) = - fanout.feeds.entry(key.data().clone()).or_insert_with(|| { - let partition = Some(key.data().clone()); - (key, PartitionFeed::new(partition, splitter.slicer), seen) - }); - let held_before = feed.held_bytes(); - let units = feed.push(part, properties)?; - fanout.held_bytes = fanout.held_bytes - held_before + feed.held_bytes(); - for unit in units { - w.write(key.clone(), unit).await.map_err(iceberg_err)?; - } - } - // Every partition stays open, so the rows they hold back share one limit. - if !properties.chooser.may_hold(fanout.held_bytes) { - for (key, feed, _) in fanout.feeds.values_mut() { - for unit in feed.release(properties)? { - w.write(key.clone(), unit).await.map_err(iceberg_err)?; - } - } - fanout.held_bytes = 0; - } - Ok(()) + InnerWriter::Fanout(fanout, splitter) => { + fanout.write(splitter, &batch, properties).await } InnerWriter::Clustered(w, splitter, live) => { for (key, part) in splitter.split_runs(&batch)? { @@ -1381,29 +1428,7 @@ impl InnerWriter { } w.close().await.map_err(iceberg_err) } - InnerWriter::Fanout(mut w, _, fanout) => { - for (_, (key, mut feed, _)) in fanout.feeds { - for unit in feed.finish(properties)? { - w.write(key.clone(), unit).await.map_err(iceberg_err)?; - } - } - let mut data_files = w.close().await.map_err(iceberg_err)?; - // `FanoutFiles` holds its per-partition writers in a `HashMap` and `close` - // iterates it directly, so the order it returns follows Rust's per-process - // `RandomState` and differs on every run. That order becomes the manifest entry - // order, then the scan-task order, then the row order of an unordered - // `SELECT *` -- which Iceberg's own - // `TestMetadataTablesWithPartitionEvolution.testPartitionColumnNamedPartition` - // compares positionally against what iceberg-java wrote - // (apache/datafusion-comet#5776). - // - // Sorting by path is what makes the task's output reproducible. Path order sorts - // by partition directory, then by the file name counter within a partition, and - // needs nothing from the write order. The clustered and unpartitioned writers - // append in creation order and are already deterministic, so only this arm sorts. - data_files.sort_unstable_by(|a, b| a.file_path().cmp(b.file_path())); - Ok(data_files) - } + InnerWriter::Fanout(fanout, _) => fanout.close(properties).await, InnerWriter::Clustered(mut w, splitter, live) => { if let Some((key, mut feed)) = live { // Feeding defers a partition's rows to here, so the unclustered-input @@ -1926,22 +1951,6 @@ impl PartitionFeed { } } -/// The feeds of a fanout write, one per partition, all open at once. -/// -/// Between them they hold back no more than a single partition may: once the rows held across all -/// partitions reach that, every partition still holding makes its choice from the rows it has. A -/// task fanning out to many partitions that each get less than a page would otherwise hold all of -/// its rows back, uncompressed, until it closed. -#[derive(Default)] -struct FanoutFeeds { - /// The `PartitionKey` is kept alongside each feed so held rows can still be written out at - /// close, and so is how many partitions the task had seen before this one, which orders - /// partitions that hold as much memory. - feeds: HashMap, - /// Bytes held back across all feeds. - held_bytes: usize, -} - /// `true` when `data_type` puts a float or double under a list or map, where a slice's offset /// window is invisible to iceberg-rust's NaN-count visitor. fn float_under_list_or_map(data_type: &DataType) -> bool { @@ -3246,11 +3255,10 @@ mod tests { /// A fanout task's data files must come back in a deterministic order. /// - /// `FanoutFiles` keeps its per-partition writers in a `HashMap` and `close` iterates it - /// directly, so the `DataFile` order it returns follows Rust's - /// per-process `RandomState` -- a different order on every run. That order becomes the - /// manifest entry order, which becomes the scan-task order, which becomes the row order - /// of an unordered `SELECT *`. Iceberg's own + /// The order a fanout write closes its files in depends on the order its rows arrived in + /// and on which partitions it closed early. That order would become the manifest entry + /// order, which becomes the scan-task order, which becomes the row order of an unordered + /// `SELECT *`. Iceberg's own /// `TestMetadataTablesWithPartitionEvolution.testPartitionColumnNamedPartition` compares /// such a `SELECT *` positionally against what iceberg-java wrote, and fails when the /// two disagree. See https://github.com/apache/datafusion-comet/issues/5776. @@ -4850,28 +4858,6 @@ mod tests { batches: Vec, writer_properties: WriterProperties, target_file_size_bytes: u64, - ) -> (TempDir, DFResult>) { - write_reporting_to( - pool, - schema, - writer_mode, - batches, - writer_properties, - target_file_size_bytes, - WriteMetrics::default(), - ) - .await - } - - /// [`write_reserving_from`], reporting to `metrics`. - async fn write_reporting_to( - pool: &Arc, - schema: Schema, - writer_mode: ProtoIcebergWriterMode, - batches: Vec, - writer_properties: WriterProperties, - target_file_size_bytes: u64, - metrics: WriteMetrics, ) -> (TempDir, DFResult>) { let temp_dir = TempDir::new().unwrap(); let spec = match writer_mode { @@ -4902,7 +4888,7 @@ mod tests { writer_properties, Some(0), Some(0), - metrics, + WriteMetrics::default(), reservation, ) .await @@ -5124,27 +5110,14 @@ mod tests { rows } - /// A fanout write whose partitions outgrow what the pool grants writes out and closes the - /// ones holding the most until what is left fits, and carries on: a closed partition's - /// next rows open a new file. Every row still lands, in more files than partitions. - #[tokio::test] - async fn a_fanout_write_the_pool_cannot_hold_closes_partitions_until_it_fits() { - // Four batches, each a unit for every one of 16 partitions, so a partition closed - // early gets more rows after it. - let batches = || { - (0..4) - .map(|batch| { - round_robin_batch_from(batch * 16 * ROWS_DIVISOR, 16 * ROWS_DIVISOR, 16) - }) - .collect::>() - }; - // One unit is one page, so every unit goes straight to its partition's file. - let properties = || { - WriterProperties::builder() - .set_data_page_row_count_limit(ROWS_DIVISOR) - .build() - }; - + /// Writes `batches` into a fanout table twice: with a pool to spare, then with half the + /// most that first write reserved, which it can only fit by closing partitions early. + /// Returns both writes' data files and the second write's directory, after checking that + /// the second fit and gave its reservation back. + async fn fanout_write_in_half_the_pool( + batches: impl Fn() -> Vec, + properties: impl Fn() -> WriterProperties, + ) -> (Vec, Vec, TempDir) { let roomy = Arc::new(PeakMemoryPool::new(usize::MAX)); let (_dir, written) = write_reserving_from( &roomy, @@ -5155,41 +5128,49 @@ mod tests { TARGET_FILE_SIZE, ) .await; - assert_eq!( - record_counts(&written.unwrap()), - vec![4 * ROWS_DIVISOR as u64; 16] - ); - + let roomy_files = written.unwrap(); let tight = Arc::new(PeakMemoryPool::new(roomy.peak() / 2)); - let metrics = WriteMetrics::default(); - let (dir, written) = write_reporting_to( + let (dir, written) = write_reserving_from( &tight, iceberg_user_schema(), ProtoIcebergWriterMode::IcebergWriterFanout, batches(), properties(), TARGET_FILE_SIZE, - metrics.clone(), ) .await; - let data_files = written.expect("closing partitions early lets the write fit"); - assert!( - metrics.files_closed_early.value() > 0, - "no partition was closed early" - ); - assert!( - data_files.len() > 16, - "{} files for 16 partitions", - data_files.len() - ); - let rows = rows_per_partition(&data_files); - assert_eq!(rows.len(), 16); - assert!( - rows.values().all(|&rows| rows == 4 * ROWS_DIVISOR as u64), - "{rows:?}" - ); - assert_eq!(parquet_files_under(dir.path()).len(), data_files.len()); + let tight_files = written.expect("closing partitions early lets the write fit"); assert_eq!(tight.reserved(), 0, "the write kept a reservation"); + (roomy_files, tight_files, dir) + } + + /// A fanout write whose partitions outgrow what the pool grants writes out and closes the + /// ones holding the most until what is left fits, and carries on: a closed partition's + /// next rows open a new file. Every row still lands, in more files than partitions. + #[tokio::test] + async fn a_fanout_write_the_pool_cannot_hold_closes_partitions_until_it_fits() { + // Four batches, each a unit for every one of 16 partitions, so a partition closed + // early gets more rows after it. One unit is one page, so every unit goes straight to + // its partition's file. + let (roomy, tight, dir) = fanout_write_in_half_the_pool( + || { + (0..4) + .map(|batch| { + round_robin_batch_from(batch * 16 * ROWS_DIVISOR, 16 * ROWS_DIVISOR, 16) + }) + .collect() + }, + || { + WriterProperties::builder() + .set_data_page_row_count_limit(ROWS_DIVISOR) + .build() + }, + ) + .await; + assert_eq!(roomy.len(), 16); + assert!(tight.len() > 16, "{} files for 16 partitions", tight.len()); + assert_eq!(rows_per_partition(&tight), rows_per_partition(&roomy)); + assert_eq!(parquet_files_under(dir.path()).len(), tight.len()); } /// A partition's memory counts the rows its feed still holds back, and closing it writes @@ -5197,40 +5178,12 @@ mod tests { /// file is open when the pool refuses: the write only fits by writing partitions out. #[tokio::test] async fn closing_a_partition_early_writes_out_the_rows_it_holds_back() { - let batches = || vec![round_robin_batch(16 * ROWS_DIVISOR, 16)]; - let properties = || WriterProperties::builder().build(); - - let roomy = Arc::new(PeakMemoryPool::new(usize::MAX)); - let (_dir, written) = write_reserving_from( - &roomy, - iceberg_user_schema(), - ProtoIcebergWriterMode::IcebergWriterFanout, - batches(), - properties(), - TARGET_FILE_SIZE, - ) - .await; - assert_eq!(written.unwrap().len(), 16); - - let tight = Arc::new(PeakMemoryPool::new(roomy.peak() / 2)); - let metrics = WriteMetrics::default(); - let (_dir, written) = write_reporting_to( - &tight, - iceberg_user_schema(), - ProtoIcebergWriterMode::IcebergWriterFanout, - batches(), - properties(), - TARGET_FILE_SIZE, - metrics.clone(), + let (roomy, tight, _dir) = fanout_write_in_half_the_pool( + || vec![round_robin_batch(16 * ROWS_DIVISOR, 16)], + || WriterProperties::builder().build(), ) .await; - let data_files = written.expect("writing held rows out lets the write fit"); - assert!( - metrics.files_closed_early.value() > 0, - "no partition was closed early" - ); - assert_eq!(record_counts(&data_files), vec![ROWS_DIVISOR as u64; 16]); - assert_eq!(tight.reserved(), 0, "the write kept a reservation"); + assert_eq!(record_counts(&tight), record_counts(&roomy)); } /// A write whose open file outgrows what the pool grants, with no partition it can close @@ -5301,7 +5254,6 @@ mod tests { let builder = MeteredParquetWriterBuilder { inner: ParquetWriterBuilder::new(WriterProperties::builder().build(), schema), open_files: open_files.clone(), - partition: None, storage: StorageWrites::Through, }; let mut closed = builder @@ -5345,7 +5297,6 @@ mod tests { Arc::clone(&schema), ), open_files: open_files.clone(), - partition: None, storage, }; let location = format!("memory:/t/{storage:?}.parquet"); diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index cf39fc47b3d..6e16b3f232d 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -3109,19 +3109,24 @@ class CometIcebergWriteActionSuite } withNativeEnabled { - var plans = Seq.empty[SparkPlan] - // withSQLConf returns Unit before Spark 4.0, so the plans are kept from inside it. + // withSQLConf returns Unit before Spark 4.0, so the assertions run inside it. withSQLConf(CometConf.COMET_OFFHEAP_MEMORY_POOL_FRACTION.key -> "0.002") { - plans = capturePlans(spark) { + val snapshot = captureWrite("fanout_oom") { spark.sql(s"INSERT INTO $catalog.$ns.fanout_oom SELECT * FROM fanout_oom_src") } + assert( + snapshot.snapshotDelta == 1L, + s"expected 1 commit, got ${snapshot.snapshotDelta}") + val writers = snapshot.plans.flatMap { plan => + collectWithSubqueries(plan) { case w: CometIcebergWriteExec => w } + } + assert( + writers.nonEmpty, + s"the write did not run natively:\n${snapshot.plans.mkString("\n--\n")}") + assert( + writers.map(_.metrics("files_closed_early").value).sum > 0, + "the write closed no partition early") } - val writers = - plans.flatMap(p => collectWithSubqueries(p) { case w: CometIcebergWriteExec => w }) - assert(writers.nonEmpty, s"the write did not run natively:\n${plans.mkString("\n--\n")}") - assert( - writers.map(_.metrics("files_closed_early").value).sum > 0, - "the write closed no partition early") // With the whole pool, every partition keeps its one file open. spark.sql(s"INSERT INTO $catalog.$ns.fanout_oom_control SELECT * FROM fanout_oom_src")