leaves12138 commented on code in PR #928:
URL: https://github.com/apache/paimon-rust/pull/928#discussion_r4083920366


##########
crates/paimon/src/arrow/format/mod.rs:
##########
@@ -137,6 +143,68 @@ pub(crate) trait FormatFileWriter: Send {
     async fn close(self: Box<Self>) -> crate::Result<FormatWriteResult>;
 }
 
+/// Account for batches retained by a format writer until its next flush.
+/// The charge is an estimate of input Arrow buffers, not encoded allocations.
+pub(crate) fn with_write_resources(
+    writer: Box<dyn FormatFileWriter>,
+    resources: Option<&crate::resource::ResourceContext>,
+) -> Box<dyn FormatFileWriter> {
+    match resources {
+        Some(resources) => Box::new(ResourceFormatWriter {
+            inner: writer,
+            reservation: resources.reservation(),
+        }),
+        None => writer,
+    }
+}
+
+struct ResourceFormatWriter {
+    inner: Box<dyn FormatFileWriter>,
+    reservation: crate::resource::MemoryReservation,
+}
+
+#[async_trait]
+impl FormatFileWriter for ResourceFormatWriter {
+    async fn write(&mut self, batch: &RecordBatch) -> crate::Result<()> {
+        self.reservation.try_grow(batch.get_array_memory_size())?;
+        self.inner.write(batch).await?;
+        if !self.inner.retains_batch_data() {
+            self.reservation.try_resize(0)?;
+        }

Review Comment:
   [P2] Release the charge for row groups that Parquet flushes internally
   
   Parquet can flush a full row group inside write() and leave a small tail in 
the next group. In that case retains_batch_data() stays true, so this code 
retains the charge for all previously flushed input batches as well. The 
reservation keeps growing across internal flushes until an entirely empty 
buffer, an explicit flush, or file close, and valid writes can fail with 
ResourceExhausted even though almost all of the charged data has already been 
flushed. This also affects the DataFusion sink, which now always enables this 
accounting.
   
   Reproduced on this head with an append table containing two Int32 columns, 
file.format=parquet, target-file-size=1gb, and write.parquet-buffer-size=64mb. 
Write the same 1,048,577-row batch three times under a 20,972,020-byte budget 
(2.5 times the batch's 8,388,808-byte Arrow estimate). Parquet's default 
1,048,576-row boundary flushes a full group on each call, but reservations 
after the first two writes are 8,388,808 and 16,777,616 bytes; the third write 
fails while requesting another 8,388,808 bytes. Without resources all three 
writes succeed, and the footer confirms row groups [1048576, 1048576, 1048576, 
3]. With resources and exactly 1,048,576 rows per batch, all three writes also 
succeed and the reservation returns to zero after each write.
   
   Please reconcile the charge with the data still retained after internal 
format flushes, releasing the emitted groups' reservation while preserving the 
unfinished group's charge, and add coverage for non-row-group-aligned batches.



-- 
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]

Reply via email to