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]