sunchao commented on code in PR #6247:
URL: https://github.com/apache/datafusion-comet/pull/6247#discussion_r4124645199


##########
native/core/src/execution/operators/iceberg_write.rs:
##########
@@ -82,11 +87,139 @@ use crate::execution::operators::iceberg_partition_path::{
 
 /// Builder chain instantiated once per task and handed to the partitioning 
wrapper.
 type IcebergDataFileWriterBuilder = DataFileWriterBuilder<
-    ParquetWriterBuilder,
+    MeteredParquetWriterBuilder,
     TrackingLocationGenerator,
     DefaultFileNameGenerator,
 >;
 
+/// 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.
+#[derive(Clone, Debug, Default)]
+struct OpenFileMemory(Arc<AtomicUsize>);
+
+impl OpenFileMemory {
+    fn bytes(&self) -> usize {
+        self.0.load(Ordering::Relaxed)
+    }
+
+    fn share(&self) -> OpenFileShare {
+        OpenFileShare {
+            total: self.clone(),
+            bytes: 0,
+        }
+    }
+}
+
+/// One open file's part of [`OpenFileMemory`], given back when the file 
closes, or when it is
+/// dropped because the task failed mid-write.
+#[derive(Debug)]
+struct OpenFileShare {
+    total: OpenFileMemory,
+    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.bytes = bytes;
+    }
+}
+
+impl Drop for OpenFileShare {
+    fn drop(&mut self) {
+        self.set(0);
+    }
+}
+
+/// [`ParquetWriterBuilder`] whose files report what they hold in memory to 
the task's
+/// [`OpenFileMemory`].
+#[derive(Clone, Debug)]
+struct MeteredParquetWriterBuilder {
+    inner: ParquetWriterBuilder,
+    open_files: OpenFileMemory,
+    /// The row-group size the files are written with, the most a file's share 
can be.
+    row_group_bytes: usize,
+}
+
+impl FileWriterBuilder for MeteredParquetWriterBuilder {
+    type R = MeteredParquetWriter;
+
+    async fn build(&self, output_file: OutputFile) -> iceberg::Result<Self::R> 
{
+        Ok(MeteredParquetWriter {
+            inner: self.inner.build(output_file).await?,
+            share: self.open_files.share(),
+            row_group_bytes: self.row_group_bytes,
+        })
+    }
+}
+
+/// iceberg-rust's [`ParquetWriter`], reporting after every write how much of 
its file it still
+/// holds in memory.
+///
+/// parquet-rs keeps a file's in-progress row group in memory and hands it to 
storage once it
+/// reaches the row-group size. `ParquetWriter` does not expose parquet-rs's 
estimate of that
+/// memory, only `current_written_size`: the bytes already handed to storage 
plus the in-progress
+/// row group's encoded size. The two agree until the file flushes its first 
row group. After
+/// that the in-progress row group is still at most one row group, so the 
share is capped at the
+/// row-group size. That over-reports a file that has flushed by at most one 
row group, which is
+/// about what it holds again just before its next flush.
+///
+/// The share does not see what parquet-rs keeps beyond the encoded estimate: 
dictionary
+/// encoders' hash tables and unencoded indices, and buffer capacity past what 
is used. Nor would
+/// it see Bloom filters, which parquet-rs sizes when a file opens and the 
eligibility gate
+/// declines today.
+struct MeteredParquetWriter {
+    inner: ParquetWriter,
+    share: OpenFileShare,
+    row_group_bytes: usize,
+}
+
+impl FileWriter for MeteredParquetWriter {
+    async fn write(&mut self, batch: &RecordBatch) -> iceberg::Result<()> {
+        self.inner.write(batch).await?;
+        self.share
+            .set(self.inner.current_written_size().min(self.row_group_bytes));

Review Comment:
   [P2] Release the row-group charge after a flush. `current_written_size()` 
includes bytes already flushed to storage, so this cap leaves every flushed 
fanout file permanently charged for a full row group until it closes. With 128 
KiB row groups, uncompressed 1,000-row batches containing distinct 256-byte 
strings, a 512 MiB file target, and a 1 MiB pool, the ninth partition requests 
1,179,648 bytes and fails with `ResourcesExhausted`, although the equivalent 
parquet-rs writers report zero live row-group bytes after every batch. Flushed 
buffers should release their reservation so subsequent partitions can reuse the 
budget. This introduces repeatable task failures for writes whose row-group 
buffers do not accumulate. Please use a live-buffer accessor from iceberg-rust 
and add a regression covering reservation reuse after flushing.
   
   Evidence: Reproduced through exact-head `run_write_task` with disposable 
test `review_flushed_fanout_files_false_oom`. The unbounded control wrote all 
12 partitions. The 1 MiB pool failed requesting another 128 KiB after reserving 
1 MiB. Parallel `ArrowWriter` controls asserted one flushed row group and 
`memory_size() == 0` for every batch. Failure cleanup deleted all files and 
returned the reservation. Reproduction patch and output are preserved at 
`/tmp/comet-6247-current-validation/reproduction.patch` and 
`/tmp/comet-6247-current-validation/task-probe.log`. Pinned iceberg-rust 
implements `current_written_size()` as `bytes_written() + in_progress_size()`.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to