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]