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 1ef7539329a..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 | FanoutWriter | ClusteredWriter +UnpartitionedWriter | FanoutPartitions | 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.** `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. @@ -281,9 +282,16 @@ 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 (`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 the storage backend from the data location's scheme and wires in Comet's S3 credential bridge when @@ -411,16 +419,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/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 diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index a04f6be61cf..60da2f74089 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 @@ -418,10 +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: 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 - 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 ba9b08bd9da..eb9b83d1f24 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -25,6 +25,7 @@ //! `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::HashMap; use std::fmt; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -44,7 +45,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 +70,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 +124,18 @@ 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 report + /// what they hold to `open_files`. + async fn build_with( + &self, + partition_key: Option, + properties: WriterProperties, + open_files: OpenFileMemory, + ) -> iceberg::Result { let parquet_builder = MeteredParquetWriterBuilder { inner: ParquetWriterBuilder::new(properties, Arc::clone(&self.schema)), - open_files: self.open_files.clone(), + open_files, storage: self.storage, }; let rolling_builder = RollingFileWriterBuilder::new( @@ -149,13 +151,28 @@ impl IcebergWriterBuilder for PartitionWriterBuilder { } } +/// How the unpartitioned and clustered writers open a partition's writer: once, with the +/// properties chosen for it. +#[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, self.open_files.clone()) + .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. +/// 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. @@ -196,16 +213,29 @@ impl PartitionProperties { /// What a task's open data files hold in memory between them. /// -/// 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 +/// 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 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 share(&self) -> OpenFileShare { @@ -226,14 +256,12 @@ struct OpenFileShare { 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); + 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; } @@ -245,7 +273,7 @@ impl Drop for OpenFileShare { } } -/// [`ParquetWriterBuilder`] whose files report what they hold in memory to the task's +/// [`ParquetWriterBuilder`] whose files report what they hold in memory to an /// [`OpenFileMemory`]. #[derive(Clone, Debug)] struct MeteredParquetWriterBuilder { @@ -816,9 +844,7 @@ 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::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}]")) @@ -843,7 +869,7 @@ impl ExecutionPlan for IcebergWriteExec { writer_properties.as_ref().clone(), partition_id, task_attempt_id, - write_time, + metrics, reservation, ) .await?; @@ -909,14 +935,35 @@ 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, +} + +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`/`FanoutWriter`/`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. The writer cannot spill, so a reservation the pool refuses 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. @@ -930,7 +977,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 @@ -1011,11 +1058,9 @@ async fn run_write_task( UnpartitionedWriter::new(data_file_builder), PartitionFeed::new(None, slicer), ), - (false, ProtoIcebergWriterMode::IcebergWriterFanout) => InnerWriter::Fanout( - FanoutWriter::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) } @@ -1033,12 +1078,15 @@ 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?; + let closed = writer + .reserve(&reservation, &open_files, &properties) + .await?; + metrics.files_closed_early.add(closed); 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 +1105,204 @@ async fn run_write_task( } } -/// Enum-based dispatch over the three iceberg-rust partitioning writers, each paired with the +/// 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::reserve`]). +struct FanoutPartitions { + builder: PartitionWriterBuilder, + /// 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, +} + +/// 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, + partitions: Vec::new(), + index: HashMap::new(), + held_bytes: 0, + closed: Vec::new(), + } + } + + /// 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() + } + + /// 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(()) + } + + /// 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() + } + + /// 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) + } +} + +/// 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 +1311,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(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( @@ -1087,17 +1328,44 @@ 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()) } } } + /// 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, + reservation: &MemoryReservation, + open_files: &OpenFileMemory, + properties: &PartitionProperties, + ) -> 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 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, + } + } + Err(refused) + } + /// 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 +1376,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)? { @@ -1116,29 +1383,8 @@ impl InnerWriter { } Ok(()) } - 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 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)? { @@ -1175,7 +1421,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)? { @@ -1183,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)?; - // `FanoutWriter` 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 @@ -1728,21 +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. - 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 { @@ -2832,7 +3040,7 @@ mod tests { writer_properties, Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await?; @@ -2932,7 +3140,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -2980,7 +3188,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -3047,11 +3255,10 @@ 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 - /// 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. @@ -4257,7 +4464,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4386,7 +4593,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4454,7 +4661,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4547,7 +4754,7 @@ mod tests { WriterProperties::builder().build(), Some(0), Some(0), - Time::default(), + WriteMetrics::default(), unbounded_reservation(), ) .await @@ -4681,7 +4888,7 @@ mod tests { writer_properties, Some(0), Some(0), - Time::default(), + WriteMetrics::default(), reservation, ) .await @@ -4894,32 +5101,127 @@ 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 + } + + /// 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, + iceberg_user_schema(), + ProtoIcebergWriterMode::IcebergWriterFanout, + batches(), + properties(), + TARGET_FILE_SIZE, + ) + .await; + let roomy_files = written.unwrap(); + let tight = Arc::new(PeakMemoryPool::new(roomy.peak() / 2)); + let (dir, written) = write_reserving_from( + &tight, + iceberg_user_schema(), + ProtoIcebergWriterMode::IcebergWriterFanout, + batches(), + properties(), + TARGET_FILE_SIZE, + ) + .await; + 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 + /// 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 (roomy, tight, _dir) = fanout_write_in_half_the_pool( + || vec![round_robin_batch(16 * ROWS_DIVISOR, 16)], + || WriterProperties::builder().build(), + ) + .await; + assert_eq!(record_counts(&tight), record_counts(&roomy)); + } + + /// 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() { - let batches = || vec![round_robin_batch(16 * ROWS_DIVISOR, 16)]; - let properties = || WriterProperties::builder().build(); + // 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, every partition's unit reaches its own file. + // 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::IcebergWriterFanout, + ProtoIcebergWriterMode::IcebergWriterUnpartitioned, batches(), properties(), TARGET_FILE_SIZE, ) .await; - assert_eq!(written.unwrap().len(), 16); - assert_eq!(parquet_files_under(dir.path()).len(), 16); + 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::IcebergWriterFanout, + ProtoIcebergWriterMode::IcebergWriterUnpartitioned, batches(), properties(), TARGET_FILE_SIZE, 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..6e16b3f232d 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,69 @@ class CometIcebergWriteActionSuite } withNativeEnabled { + // withSQLConf returns Unit before Spark 4.0, so the assertions run inside it. withSQLConf(CometConf.COMET_OFFHEAP_MEMORY_POOL_FRACTION.key -> "0.002") { - val (failedPlans, error) = captureFailedPlans(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") + } + + // 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") + // 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) + } + } + } + + // 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 +3184,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) } } }