This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new 9978b9dc [rust] Account writer memory with ResourceContext (#928)
9978b9dc is described below

commit 9978b9dc2ca89171d21ead8df38c09e2bda52eba
Author: Jingsong Lee <[email protected]>
AuthorDate: Wed Sep 23 23:23:33 2026 +0800

    [rust] Account writer memory with ResourceContext (#928)
---
 crates/integrations/datafusion/src/memory.rs       |   9 +
 .../datafusion/src/physical_plan/sink.rs           |   4 +-
 crates/paimon/src/arrow/format/mod.rs              | 145 ++++++++-
 crates/paimon/src/arrow/format/parquet.rs          |  69 +++-
 crates/paimon/src/arrow/format/shredding.rs        |  57 ++++
 crates/paimon/src/resource/memory.rs               |   2 +-
 crates/paimon/src/resource/mod.rs                  |   7 +-
 crates/paimon/src/table/data_file_writer.rs        |  20 +-
 .../src/table/dedicated_format_file_writer.rs      |  12 +
 crates/paimon/src/table/format_write_builder.rs    |   5 +
 crates/paimon/src/table/kv_file_writer.rs          |  30 +-
 crates/paimon/src/table/postpone_file_writer.rs    |  16 +-
 .../src/table/postpone_fixed_bucket_write.rs       |   7 +-
 .../table/postpone_fixed_bucket_write_builder.rs   |   9 +
 crates/paimon/src/table/table_write.rs             | 222 ++++++++-----
 crates/paimon/src/table/write_builder.rs           |  25 +-
 crates/paimon/tests/writer_resources_test.rs       | 359 +++++++++++++++++++++
 17 files changed, 898 insertions(+), 100 deletions(-)

diff --git a/crates/integrations/datafusion/src/memory.rs 
b/crates/integrations/datafusion/src/memory.rs
index 7419b218..bd16c3a2 100644
--- a/crates/integrations/datafusion/src/memory.rs
+++ b/crates/integrations/datafusion/src/memory.rs
@@ -55,3 +55,12 @@ pub(crate) fn reader_resources(
         .memory_pool(Arc::new(DataFusionMemoryPool { reservation }))
         .build()
 }
+
+pub(crate) fn writer_resources(context: &TaskContext) -> 
paimon::Result<ResourceContext> {
+    let reservation = MemoryConsumer::new("PaimonTableWrite")
+        .with_can_spill(false)
+        .register(context.memory_pool());
+    ResourceContext::builder()
+        .memory_pool(Arc::new(DataFusionMemoryPool { reservation }))
+        .build()
+}
diff --git a/crates/integrations/datafusion/src/physical_plan/sink.rs 
b/crates/integrations/datafusion/src/physical_plan/sink.rs
index dd07a9ab..72957c5d 100644
--- a/crates/integrations/datafusion/src/physical_plan/sink.rs
+++ b/crates/integrations/datafusion/src/physical_plan/sink.rs
@@ -137,13 +137,15 @@ impl DataSink for PaimonDataSink {
     async fn write_all(
         &self,
         mut data: SendableRecordBatchStream,
-        _context: &Arc<TaskContext>,
+        context: &Arc<TaskContext>,
     ) -> DFResult<u64> {
         let wb = if self.overwrite {
             self.table.new_write_builder().with_overwrite()
         } else {
             self.table.new_write_builder()
         };
+        let resources = 
crate::memory::writer_resources(context).map_err(to_datafusion_error)?;
+        let wb = wb.with_resources(resources);
         let mut tw = wb.new_write().map_err(to_datafusion_error)?;
         let mut row_count = 0u64;
 
diff --git a/crates/paimon/src/arrow/format/mod.rs 
b/crates/paimon/src/arrow/format/mod.rs
index e9c057e3..8f90214f 100644
--- a/crates/paimon/src/arrow/format/mod.rs
+++ b/crates/paimon/src/arrow/format/mod.rs
@@ -39,7 +39,7 @@ use crate::Error;
 use arrow_array::RecordBatch;
 use arrow_schema::SchemaRef;
 use async_trait::async_trait;
-use std::collections::HashMap;
+use std::collections::{HashMap, VecDeque};
 use std::sync::Arc;
 
 /// Predicates with the file-level field context needed for pushdown.
@@ -114,6 +114,19 @@ pub(crate) trait FormatFileWriter: Send {
     /// Number of bytes buffered in the current row group (not yet flushed).
     fn in_progress_size(&self) -> usize;
 
+    /// Whether this writer still retains batch data. This can differ from
+    /// `in_progress_size` while a format is inferring its physical schema.
+    fn retains_batch_data(&self) -> bool {
+        self.in_progress_size() != 0
+    }
+
+    /// Number of input rows still represented by buffered data, when known.
+    /// Unlike `in_progress_size`, this allows accounting to release the part
+    /// of a batch already written by an automatic row-group flush.
+    fn pending_rows(&self) -> Option<usize> {
+        None
+    }
+
     /// Flush the current row group to storage without closing the file.
     async fn flush(&mut self) -> crate::Result<()>;
 
@@ -137,6 +150,136 @@ 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(),
+            batches: VecDeque::new(),
+            charged_rows: 0,
+        }),
+        None => writer,
+    }
+}
+
+struct ResourceFormatWriter {
+    inner: Box<dyn FormatFileWriter>,
+    reservation: crate::resource::MemoryReservation,
+    batches: VecDeque<BufferedBatchCharge>,
+    charged_rows: usize,
+}
+
+struct BufferedBatchCharge {
+    rows: usize,
+    bytes: usize,
+}
+
+impl ResourceFormatWriter {
+    fn release_flushed(&mut self) -> crate::Result<()> {
+        let Some(pending_rows) = self.inner.pending_rows() else {
+            if !self.inner.retains_batch_data() {
+                self.batches.clear();
+                self.charged_rows = 0;
+                self.reservation.try_resize(0)?;
+            }
+            return Ok(());
+        };
+        if pending_rows == 0 {
+            self.batches.clear();
+            self.charged_rows = 0;
+            return self.reservation.try_resize(0);
+        }
+        if pending_rows > self.charged_rows {
+            // Keep the full charge if a writer reports more rows than we have
+            // observed; releasing any amount would risk under-accounting.
+            return Ok(());
+        }
+        let mut flushed_rows = self.charged_rows - pending_rows;
+        let mut released_bytes = 0;
+        while flushed_rows > 0 {
+            let batch = self.batches.front_mut().expect("charged rows remain");
+            let consumed_rows = flushed_rows.min(batch.rows);
+            let remaining_rows = batch.rows - consumed_rows;
+            let remaining_bytes = (batch.bytes as u128 * remaining_rows as 
u128)
+                .div_ceil(batch.rows as u128) as usize;
+            released_bytes += batch.bytes - remaining_bytes;
+            flushed_rows -= consumed_rows;
+            if remaining_rows == 0 {
+                self.batches.pop_front();
+            } else {
+                batch.rows = remaining_rows;
+                batch.bytes = remaining_bytes;
+            }
+        }
+        self.charged_rows = pending_rows;
+        self.reservation
+            .try_resize(self.reservation.size() - released_bytes)
+    }
+}
+
+#[async_trait]
+impl FormatFileWriter for ResourceFormatWriter {
+    async fn write(&mut self, batch: &RecordBatch) -> crate::Result<()> {
+        let bytes = if batch.num_rows() == 0 {
+            0
+        } else {
+            batch.get_array_memory_size()
+        };
+        self.reservation.try_grow(bytes)?;
+        self.inner.write(batch).await?;
+        if batch.num_rows() != 0 {
+            self.batches.push_back(BufferedBatchCharge {
+                rows: batch.num_rows(),
+                bytes,
+            });
+            self.charged_rows += batch.num_rows();
+        }
+        self.release_flushed()
+    }
+
+    fn num_bytes(&self) -> usize {
+        self.inner.num_bytes()
+    }
+
+    fn in_progress_size(&self) -> usize {
+        self.inner.in_progress_size()
+    }
+
+    fn retains_batch_data(&self) -> bool {
+        self.inner.retains_batch_data()
+    }
+
+    fn pending_rows(&self) -> Option<usize> {
+        self.inner.pending_rows()
+    }
+
+    async fn flush(&mut self) -> crate::Result<()> {
+        self.inner.flush().await?;
+        self.release_flushed()
+    }
+
+    fn commit_field_metadata(
+        &mut self,
+        metadata: &crate::arrow::shredding::FieldMetadata,
+    ) -> crate::Result<()> {
+        self.inner.commit_field_metadata(metadata)
+    }
+
+    async fn close(self: Box<Self>) -> crate::Result<FormatWriteResult> {
+        let Self {
+            inner, reservation, ..
+        } = *self;
+        let result = inner.close().await;
+        drop(reservation);
+        result
+    }
+}
+
 pub(crate) struct FormatWriteResult {
     pub(crate) file_size: u64,
     pub(crate) value_stats: Option<FormatValueStats>,
diff --git a/crates/paimon/src/arrow/format/parquet.rs 
b/crates/paimon/src/arrow/format/parquet.rs
index 22d02ce8..2d203bba 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -465,6 +465,10 @@ impl FormatFileWriter for ParquetFormatWriter {
         self.inner.in_progress_size()
     }
 
+    fn pending_rows(&self) -> Option<usize> {
+        Some(self.inner.in_progress_rows())
+    }
+
     async fn flush(&mut self) -> crate::Result<()> {
         self.inner
             .flush()
@@ -2748,10 +2752,12 @@ mod tests {
         PredicateOperator, RowSelection,
     };
     use crate::arrow::format::{
-        create_format_reader, create_format_writer, FormatFileReader, 
FormatFileWriter,
+        create_format_reader, create_format_writer, with_write_resources, 
FormatFileReader,
+        FormatFileWriter,
     };
     use crate::arrow::{build_target_arrow_schema, variant_arrow_type, 
ReadBudget};
     use crate::io::FileIOBuilder;
+    use crate::resource::ResourceContext;
     use crate::spec::{
         ArrayType, BigIntType, DataField, DataType, Datum, IntType, 
LocalZonedTimestampType,
         MapType, PredicateBuilder, TimestampType, VarCharType, VariantType,
@@ -2767,7 +2773,7 @@ mod tests {
     use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema 
as ArrowSchema};
     use futures::{StreamExt, TryStreamExt};
     use parquet::basic::{Compression, GzipLevel, ZstdLevel};
-    use parquet::file::properties::EnabledStatistics;
+    use parquet::file::properties::{EnabledStatistics, WriterProperties};
     use parquet::file::statistics::Statistics as ParquetStatistics;
     use parquet::schema::{parser::parse_message_type, types::SchemaDescriptor};
     use std::collections::HashMap;
@@ -3662,6 +3668,65 @@ mod tests {
             .unwrap();
     }
 
+    #[tokio::test]
+    async fn test_resource_charge_releases_auto_flushed_row_group_with_tail() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/test_parquet_writer_resource_auto_flush.parquet";
+        let output = file_io.new_output(path).unwrap();
+        let schema = writer_arrow_schema();
+        let props = WriterProperties::builder()
+            .set_max_row_group_row_count(Some(4))
+            .build();
+        let inner = AsyncArrowWriter::try_new(
+            output.async_writer().await.unwrap(),
+            schema.clone(),
+            Some(props),
+        )
+        .unwrap();
+        let first = writer_test_batch(&schema, vec![1, 2, 3, 4, 5], vec![10, 
20, 30, 40, 50]);
+        let second = writer_test_batch(&schema, vec![6, 7, 8, 9], vec![60, 70, 
80, 90]);
+        let first_bytes = first.get_array_memory_size();
+        let second_bytes = second.get_array_memory_size();
+        let resources = ResourceContext::builder()
+            .memory_limit(first_bytes + second_bytes / 2)
+            .build()
+            .unwrap();
+        let mut writer = with_write_resources(
+            Box::new(ParquetFormatWriter {
+                inner,
+                input_schema: schema.clone(),
+                schema,
+                write_fields: None,
+                stats_modes: None,
+                stats_dense_store: false,
+            }),
+            Some(&resources),
+        );
+
+        writer.write(&first).await.unwrap();
+        assert!(writer.in_progress_size() > 0);
+        assert_eq!(
+            resources.metrics().reserved_memory_bytes,
+            first_bytes.div_ceil(5)
+        );
+        writer.write(&second).await.unwrap();
+        assert_eq!(
+            resources.metrics().reserved_memory_bytes,
+            second_bytes.div_ceil(4)
+        );
+        writer.close().await.unwrap();
+        assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+
+        let bytes = file_io.new_input(path).unwrap().read().await.unwrap();
+        let reader =
+            
parquet::arrow::arrow_reader::ParquetRecordBatchReader::try_new(bytes, 
1024).unwrap();
+        let total_rows: usize = reader
+            .into_iter()
+            .map(|batch| batch.unwrap().num_rows())
+            .sum();
+        assert_eq!(total_rows, 9);
+    }
+
     #[tokio::test]
     async fn test_parquet_writer_write_and_close() {
         let file_io = FileIOBuilder::new("memory").build().unwrap();
diff --git a/crates/paimon/src/arrow/format/shredding.rs 
b/crates/paimon/src/arrow/format/shredding.rs
index f0446e4c..7ebe3519 100644
--- a/crates/paimon/src/arrow/format/shredding.rs
+++ b/crates/paimon/src/arrow/format/shredding.rs
@@ -350,6 +350,26 @@ impl FormatFileWriter for ShreddingFormatWriter {
         }
     }
 
+    fn retains_batch_data(&self) -> bool {
+        match &self.state {
+            ShreddingWriterState::Ready { inner, .. } => 
inner.retains_batch_data(),
+            ShreddingWriterState::Infer {
+                buffered_batches, ..
+            } => !buffered_batches.is_empty(),
+            ShreddingWriterState::Closed => false,
+        }
+    }
+
+    fn pending_rows(&self) -> Option<usize> {
+        match &self.state {
+            ShreddingWriterState::Ready { inner, .. } => inner.pending_rows(),
+            ShreddingWriterState::Infer {
+                buffered_row_count, ..
+            } => Some(*buffered_row_count),
+            ShreddingWriterState::Closed => Some(0),
+        }
+    }
+
     async fn flush(&mut self) -> crate::Result<()> {
         self.finalize_inferred_writer().await?;
         match &mut self.state {
@@ -384,7 +404,12 @@ impl FormatFileWriter for ShreddingFormatWriter {
 mod tests {
     use super::*;
     use crate::arrow::build_target_arrow_schema;
+    use crate::arrow::format::with_write_resources;
+    use crate::resource::ResourceContext;
     use crate::spec::{DataType, IntType, MapType, VarCharType, VariantType};
+    use arrow_array::Int32Array;
+    use arrow_schema::{DataType as ArrowDataType, Field, Schema};
+    use std::sync::Arc;
 
     /// Factory that must never be reached: these tests only exercise plan
     /// detection, which fails before any writer is created.
@@ -412,6 +437,38 @@ mod tests {
         )
     }
 
+    #[tokio::test]
+    async fn inference_buffer_is_charged_without_triggering_row_group_flush() {
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "value",
+            ArrowDataType::Int32,
+            false,
+        )]));
+        let batch =
+            RecordBatch::try_new(schema.clone(), 
vec![Arc::new(Int32Array::from(vec![1, 2]))])
+                .unwrap();
+        let writer = ShreddingFormatWriter {
+            state: ShreddingWriterState::Infer {
+                writer_factory: Some(Box::new(NoopWriterFactory)),
+                schema,
+                logical_write_fields: vec![],
+                format_options: HashMap::new(),
+                buffered_batches: vec![],
+                buffered_row_count: 0,
+                infer_buffer_row_count: 10,
+                plan_builder: InferPlanBuilder::Variant,
+            },
+            compression: "zstd".to_string(),
+        };
+        let resources = ResourceContext::builder().build().unwrap();
+        let mut writer = with_write_resources(Box::new(writer), 
Some(&resources));
+        writer.write(&batch).await.unwrap();
+        assert_eq!(writer.in_progress_size(), 0);
+        assert!(resources.metrics().reserved_memory_bytes > 0);
+        drop(writer);
+        assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    }
+
     /// Mirroring Java's `ShreddingWritePlanWriterFactories`: at most one
     /// shredding plan may be active for a file.
     #[tokio::test]
diff --git a/crates/paimon/src/resource/memory.rs 
b/crates/paimon/src/resource/memory.rs
index 4be2feea..05f904ed 100644
--- a/crates/paimon/src/resource/memory.rs
+++ b/crates/paimon/src/resource/memory.rs
@@ -25,7 +25,7 @@ use crate::{Error, Result};
 /// Calls may run concurrently on any runtime thread. Failed reservations must
 /// leave the pool unchanged. Methods must not panic; `release` is infallible 
and
 /// may run during unwinding. The pool must not wait for, or call back into, a
-/// reader to free memory. Reclamation belongs to the consumer, outside the 
pool.
+/// consumer to free memory. Reclamation belongs to the consumer, outside the 
pool.
 pub trait MemoryPool: std::fmt::Debug + Send + Sync + 'static {
     fn try_reserve(&self, bytes: usize) -> Result<()>;
     fn release(&self, bytes: usize);
diff --git a/crates/paimon/src/resource/mod.rs 
b/crates/paimon/src/resource/mod.rs
index 24e99f18..36245939 100644
--- a/crates/paimon/src/resource/mod.rs
+++ b/crates/paimon/src/resource/mod.rs
@@ -15,7 +15,7 @@
 // specific language governing permissions and limitations
 // under the License.
 
-//! Shared memory reservations for readers and embedding engines.
+//! Shared memory reservations for readers, writers, and embedding engines.
 //!
 //! Parquet readers reserve each selected row group's projected uncompressed
 //! column size before data I/O, including columns needed by decoder 
predicates.
@@ -29,6 +29,9 @@
 //! Reservations are accounting, not an allocator or an RSS limit. Row-group
 //! charges are estimates; metadata, merge state, transient batches and 
allocation
 //! overhead are not fully covered, and actual memory can exceed the estimates.
+//! Writers charge retained key-value input batches and unflushed format-writer
+//! input batches. Sorting, encoding, transient routing batches, and file 
indexes
+//! are not fully covered. Input batches held by the caller remain its 
responsibility.
 
 mod memory;
 pub use memory::{MemoryPool, MemoryReservation, ResourceMetrics};
@@ -50,7 +53,7 @@ use crate::Result;
 /// let resources = ResourceContext::builder()
 ///     .memory_limit(256 * 1024 * 1024)
 ///     .build()?;
-/// // Pass resources.clone() to ReadBuilder::with_resources.
+/// // Pass resources.clone() to ReadBuilder::with_resources or 
WriteBuilder::with_resources.
 /// // Other consumers reserve from the same budget for their own retained 
state.
 /// let mut reservation = resources.reservation();
 /// reservation.try_grow(1024)?;
diff --git a/crates/paimon/src/table/data_file_writer.rs 
b/crates/paimon/src/table/data_file_writer.rs
index 2ace11e2..b25c4a04 100644
--- a/crates/paimon/src/table/data_file_writer.rs
+++ b/crates/paimon/src/table/data_file_writer.rs
@@ -23,8 +23,11 @@
 //! [`DataFileMeta`] for the commit path.
 
 use super::data_file_index_writer::{DataFileIndexWriter, FileIndexOptions};
-use crate::arrow::format::{create_format_writer, FormatFileWriter, 
FormatValueStats};
+use crate::arrow::format::{
+    create_format_writer, with_write_resources, FormatFileWriter, 
FormatValueStats,
+};
 use crate::io::FileIO;
+use crate::resource::ResourceContext;
 use crate::spec::data_file_to_file_index_file_name;
 use crate::spec::stats::BinaryTableStats;
 use crate::spec::{bucket_path_under, DataField, DataFileMeta, 
EMPTY_SERIALIZED_ROW};
@@ -67,6 +70,7 @@ pub(crate) struct DataFileWriter {
     current_row_count: i64,
     index_options: Option<Arc<FileIndexOptions>>,
     current_index: Option<DataFileIndexWriter>,
+    resources: Option<ResourceContext>,
     /// Paths owned by this write until prepare_commit hands them to the 
caller.
     created_paths: Vec<String>,
 }
@@ -113,6 +117,7 @@ impl DataFileWriter {
             current_row_count: 0,
             index_options: None,
             current_index: None,
+            resources: None,
             created_paths: Vec::new(),
         }
     }
@@ -122,10 +127,19 @@ impl DataFileWriter {
         self
     }
 
+    pub(crate) fn with_resources(mut self, resources: Option<ResourceContext>) 
-> Self {
+        self.resources = resources;
+        self
+    }
+
+    pub(crate) fn set_resources(&mut self, resources: Option<ResourceContext>) 
{
+        self.resources = resources;
+    }
+
     /// Write a RecordBatch. Rolls to a new file when target size is reached.
     pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
         let result = self.write_batch(batch).await;
-        if result.is_err() && self.index_options.is_some() {
+        if self.index_options.is_some() && result.is_err() {
             self.abort().await;
         }
         result
@@ -195,7 +209,7 @@ impl DataFileWriter {
             Some(&self.format_options),
         )
         .await?;
-        self.current_writer = Some(writer);
+        self.current_writer = Some(with_write_resources(writer, 
self.resources.as_ref()));
         self.current_index = index;
         self.current_file_name = Some(file_name);
         self.current_row_count = 0;
diff --git a/crates/paimon/src/table/dedicated_format_file_writer.rs 
b/crates/paimon/src/table/dedicated_format_file_writer.rs
index b9fa2714..d65b8744 100644
--- a/crates/paimon/src/table/dedicated_format_file_writer.rs
+++ b/crates/paimon/src/table/dedicated_format_file_writer.rs
@@ -16,6 +16,7 @@
 // under the License.
 
 use crate::io::FileIO;
+use crate::resource::ResourceContext;
 use crate::spec::{BlobViewStruct, DataField, DataFileMeta, DataType};
 use crate::table::data_file_writer::DataFileWriter;
 use crate::Result;
@@ -200,6 +201,17 @@ impl AppendDedicatedFormatFileWriter {
         }
     }
 
+    pub(crate) fn with_resources(mut self, resources: Option<ResourceContext>) 
-> Self {
+        self.normal_writer.set_resources(resources.clone());
+        for blob in &mut self.blob_writers {
+            blob.writer.set_resources(resources.clone());
+        }
+        if let Some(vector) = &mut self.vector_writer {
+            vector.writer.set_resources(resources);
+        }
+        self
+    }
+
     pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
         if batch.num_rows() == 0 {
             return Ok(());
diff --git a/crates/paimon/src/table/format_write_builder.rs 
b/crates/paimon/src/table/format_write_builder.rs
index ffd12af8..fc1764c1 100644
--- a/crates/paimon/src/table/format_write_builder.rs
+++ b/crates/paimon/src/table/format_write_builder.rs
@@ -19,6 +19,7 @@
 
 use super::write_builder::validate_commit_user;
 use super::{DataEvolutionDeleteWriter, Table, TableCommit, TableUpdate, 
TableWrite};
+use crate::resource::ResourceContext;
 use uuid::Uuid;
 
 pub(crate) struct FormatWriteBuilder<'a> {
@@ -52,6 +53,10 @@ impl<'a> FormatWriteBuilder<'a> {
         self
     }
 
+    pub(crate) fn with_resources(self, _resources: ResourceContext) -> Self {
+        self
+    }
+
     pub(crate) fn new_commit(&self) -> TableCommit {
         TableCommit::new(self.table.clone(), self.commit_user.clone())
     }
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 24f3e3af..8cc7efbf 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -27,8 +27,9 @@
 //! Reference: 
[org.apache.paimon.io.KeyValueDataFileWriterImpl](https://github.com/apache/paimon/blob/release-1.3/paimon-core/src/main/java/org/apache/paimon/io/KeyValueDataFileWriterImpl.java)
 
 use crate::arrow::arrow_fields_to_paimon;
-use crate::arrow::format::create_format_writer;
+use crate::arrow::format::{create_format_writer, with_write_resources};
 use crate::io::FileIO;
+use crate::resource::{MemoryReservation, ResourceContext};
 use crate::spec::stats::{compute_column_stats, BinaryTableStats};
 use crate::spec::{
     bucket_path_under, extract_datum_from_arrow, AggregationConfig, 
BinaryRowBuilder, CoreOptions,
@@ -58,6 +59,8 @@ pub(crate) struct KeyValueFileWriter {
     buffer: Vec<RecordBatch>,
     /// Approximate buffered bytes.
     buffer_bytes: usize,
+    resources: Option<ResourceContext>,
+    buffer_reservation: Option<MemoryReservation>,
     /// Completed file metadata.
     written_files: Vec<DataFileMeta>,
     /// Completed changelog file metadata.
@@ -147,11 +150,19 @@ impl KeyValueFileWriter {
             next_sequence_number,
             buffer: Vec::new(),
             buffer_bytes: 0,
+            resources: None,
+            buffer_reservation: None,
             written_files: Vec::new(),
             written_changelog_files: Vec::new(),
         })
     }
 
+    pub(crate) fn with_resources(mut self, resources: Option<ResourceContext>) 
-> Self {
+        self.buffer_reservation = 
resources.as_ref().map(ResourceContext::reservation);
+        self.resources = resources;
+        self
+    }
+
     /// Buffer a RecordBatch. Flushes when buffer exceeds write_buffer_size.
     /// Sequence numbers are assigned per-bucket on flush, matching Java 
Paimon behavior.
     pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
@@ -170,6 +181,9 @@ impl KeyValueFileWriter {
             .iter()
             .map(|c| c.get_buffer_memory_size())
             .sum();
+        if let Some(reservation) = &mut self.buffer_reservation {
+            reservation.try_grow(batch_bytes)?;
+        }
         self.buffer.push(batch);
         self.buffer_bytes += batch_bytes;
 
@@ -194,6 +208,8 @@ impl KeyValueFileWriter {
 
         let batches = std::mem::take(&mut self.buffer);
         self.buffer_bytes = 0;
+        let _buffer_reservation = self.buffer_reservation.take();
+        self.buffer_reservation = 
self.resources.as_ref().map(ResourceContext::reservation);
 
         // Concatenate all buffered batches, then immediately free the 
originals.
         let user_schema = batches[0].schema();
@@ -416,7 +432,7 @@ impl KeyValueFileWriter {
         self.file_io.mkdirs(&format!("{bucket_dir}/")).await?;
         let file_path = format!("{bucket_dir}/{file_name}");
         let output = self.file_io.new_output(&file_path)?;
-        let mut writer = create_format_writer(
+        let writer = create_format_writer(
             &output,
             physical_schema.clone(),
             write.file_compression,
@@ -426,6 +442,7 @@ impl KeyValueFileWriter {
             None,
         )
         .await?;
+        let mut writer = with_write_resources(writer, self.resources.as_ref());
 
         let vk_idx = batch
             .schema()
@@ -484,7 +501,11 @@ impl KeyValueFileWriter {
                     message: format!("Failed to create physical batch: {e}"),
                     source: None,
                 })?;
-            writer.write(&chunk_batch).await?;
+            if let Err(error) = writer.write(&chunk_batch).await {
+                let _ = writer.close().await;
+                let _ = self.file_io.delete_file(&file_path).await;
+                return Err(error);
+            }
         }
 
         let file_size = writer.close().await?.file_size as i64;
@@ -893,6 +914,9 @@ impl KeyValueFileWriter {
     pub(crate) async fn abort(&mut self) {
         self.buffer.clear();
         self.buffer_bytes = 0;
+        if let Some(reservation) = &mut self.buffer_reservation {
+            let _ = reservation.try_resize(0);
+        }
         let bucket_path = bucket_path_under(
             &self.config.table_location,
             &self.config.partition_path,
diff --git a/crates/paimon/src/table/postpone_file_writer.rs 
b/crates/paimon/src/table/postpone_file_writer.rs
index 1e3b64a4..c683a99c 100644
--- a/crates/paimon/src/table/postpone_file_writer.rs
+++ b/crates/paimon/src/table/postpone_file_writer.rs
@@ -24,8 +24,9 @@
 //!
 //! Reference: 
[PostponeBucketWriter](https://github.com/apache/paimon/blob/release-1.3/paimon-core/src/main/java/org/apache/paimon/table/sink/PostponeBucketWriter.java)
 
-use crate::arrow::format::{create_format_writer, FormatFileWriter};
+use crate::arrow::format::{create_format_writer, with_write_resources, 
FormatFileWriter};
 use crate::io::FileIO;
+use crate::resource::ResourceContext;
 use crate::spec::stats::BinaryTableStats;
 use crate::spec::{bucket_path_under, DataFileMeta, EMPTY_SERIALIZED_ROW, 
VALUE_KIND_FIELD_NAME};
 use crate::table::kv_file_writer::build_physical_schema;
@@ -70,6 +71,7 @@ pub(crate) struct PostponeFileWriter {
     created_paths: Vec<String>,
     /// Background file close tasks spawned during rolling.
     in_flight_closes: JoinSet<Result<DataFileMeta>>,
+    resources: Option<ResourceContext>,
 }
 
 impl PostponeFileWriter {
@@ -86,9 +88,15 @@ impl PostponeFileWriter {
             written_files: Vec::new(),
             created_paths: Vec::new(),
             in_flight_closes: JoinSet::new(),
+            resources: None,
         }
     }
 
+    pub(crate) fn with_resources(mut self, resources: Option<ResourceContext>) 
-> Self {
+        self.resources = resources;
+        self
+    }
+
     pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
         if batch.num_rows() == 0 {
             return Ok(());
@@ -101,7 +109,6 @@ impl PostponeFileWriter {
         let num_rows = batch.num_rows();
         let start_seq = self.next_sequence_number;
         let end_seq = start_seq + num_rows as i64 - 1;
-        self.next_sequence_number = end_seq + 1;
 
         // Build physical batch: [_SEQUENCE_NUMBER, _VALUE_KIND, 
all_user_cols...]
         let mut physical_columns: Vec<Arc<dyn arrow_array::Array>> = 
Vec::new();
@@ -134,12 +141,13 @@ impl PostponeFileWriter {
                 }
             })?;
 
-        self.current_row_count += num_rows as i64;
         self.current_writer
             .as_mut()
             .unwrap()
             .write(&physical_batch)
             .await?;
+        self.next_sequence_number = end_seq + 1;
+        self.current_row_count += num_rows as i64;
 
         // Roll to a new file if target size is reached — close in background
         if self.current_writer.as_ref().unwrap().num_bytes() as i64 >= 
self.config.target_file_size
@@ -240,7 +248,7 @@ impl PostponeFileWriter {
             None,
         )
         .await?;
-        self.current_writer = Some(writer);
+        self.current_writer = Some(with_write_resources(writer, 
self.resources.as_ref()));
         self.current_file_name = Some(file_name);
         self.current_row_count = 0;
         self.current_file_start_seq = self.next_sequence_number;
diff --git a/crates/paimon/src/table/postpone_fixed_bucket_write.rs 
b/crates/paimon/src/table/postpone_fixed_bucket_write.rs
index 488e3a3f..beb20bcc 100644
--- a/crates/paimon/src/table/postpone_fixed_bucket_write.rs
+++ b/crates/paimon/src/table/postpone_fixed_bucket_write.rs
@@ -19,6 +19,7 @@ use super::postpone_bucket_plan::data_invalid;
 use super::postpone_fixed_bucket_router::{
     validate_postpone_fixed_bucket_table, PostponeFixedBucketRouter,
 };
+use crate::resource::ResourceContext;
 use crate::spec::CoreOptions;
 use crate::table::{CommitMessage, PostponeBucketPlan, Table, TableCommit, 
TableWrite};
 use crate::Result;
@@ -38,9 +39,13 @@ impl PostponeFixedBucketTableWrite {
         commit_user: String,
         plan: PostponeBucketPlan,
         overwrite: bool,
+        resources: Option<ResourceContext>,
     ) -> Result<Self> {
         validate_postpone_fixed_bucket_write(table)?;
-        let inner = TableWrite::new(table, commit_user)?;
+        let mut inner = TableWrite::new(table, commit_user)?;
+        if let Some(resources) = resources {
+            inner = inner.with_resources(resources);
+        }
         Ok(Self {
             inner: if overwrite {
                 inner.with_overwrite()
diff --git a/crates/paimon/src/table/postpone_fixed_bucket_write_builder.rs 
b/crates/paimon/src/table/postpone_fixed_bucket_write_builder.rs
index 44905876..3439e155 100644
--- a/crates/paimon/src/table/postpone_fixed_bucket_write_builder.rs
+++ b/crates/paimon/src/table/postpone_fixed_bucket_write_builder.rs
@@ -17,6 +17,7 @@
 
 use super::postpone_bucket_plan::data_invalid;
 use super::postpone_fixed_bucket_router::validate_postpone_fixed_bucket_table;
+use crate::resource::ResourceContext;
 use crate::table::write_builder::{ensure_table_write_allowed, 
validate_commit_user};
 use crate::table::{
     PostponeBucketPlan, PostponeFixedBucketTableCommit, 
PostponeFixedBucketTableWrite, Table,
@@ -29,6 +30,7 @@ pub struct PostponeFixedBucketWriteBuilder<'a> {
     commit_user: String,
     overwrite: bool,
     bucket_plan: Option<PostponeBucketPlan>,
+    resources: Option<ResourceContext>,
 }
 
 impl<'a> PostponeFixedBucketWriteBuilder<'a> {
@@ -39,6 +41,7 @@ impl<'a> PostponeFixedBucketWriteBuilder<'a> {
             commit_user: Uuid::new_v4().to_string(),
             overwrite: false,
             bucket_plan: None,
+            resources: None,
         })
     }
 
@@ -63,6 +66,11 @@ impl<'a> PostponeFixedBucketWriteBuilder<'a> {
         self
     }
 
+    pub fn with_resources(mut self, resources: ResourceContext) -> Self {
+        self.resources = Some(resources);
+        self
+    }
+
     pub fn new_commit(&self) -> PostponeFixedBucketTableCommit {
         PostponeFixedBucketTableCommit::new(self.table, 
self.commit_user.clone(), self.overwrite)
     }
@@ -83,6 +91,7 @@ impl<'a> PostponeFixedBucketWriteBuilder<'a> {
             self.commit_user.clone(),
             plan,
             self.overwrite,
+            self.resources.clone(),
         )
     }
 }
diff --git a/crates/paimon/src/table/table_write.rs 
b/crates/paimon/src/table/table_write.rs
index fa819f13..77b186c4 100644
--- a/crates/paimon/src/table/table_write.rs
+++ b/crates/paimon/src/table/table_write.rs
@@ -22,6 +22,7 @@
 
 use crate::arrow::build_target_arrow_schema;
 use crate::arrow::partition::partition_array;
+use crate::resource::ResourceContext;
 use crate::spec::PartitionComputer;
 use crate::spec::{
     first_row_supports_changelog_producer, BinaryRow, ChangelogProducer, 
CoreOptions, DataField,
@@ -89,13 +90,24 @@ impl FileWriter {
     }
 
     async fn prepare_commit(mut self) -> Result<PreparedFiles> {
+        let result = match &mut self {
+            FileWriter::Append(w) => 
w.prepare_commit().await.map(PreparedFiles::data),
+            FileWriter::AppendDedicated(w) => 
w.prepare_commit().await.map(PreparedFiles::data),
+            FileWriter::KeyValue(w) => w.prepare_commit().await,
+            FileWriter::Postpone(w) => 
w.prepare_commit().await.map(PreparedFiles::data),
+        };
+        if result.is_err() {
+            self.abort().await;
+        }
+        result
+    }
+
+    async fn abort(&mut self) {
         match self {
-            FileWriter::Append(ref mut w) => 
w.prepare_commit().await.map(PreparedFiles::data),
-            FileWriter::AppendDedicated(ref mut w) => {
-                w.prepare_commit().await.map(PreparedFiles::data)
-            }
-            FileWriter::KeyValue(ref mut w) => w.prepare_commit().await,
-            FileWriter::Postpone(ref mut w) => 
w.prepare_commit().await.map(PreparedFiles::data),
+            FileWriter::Append(w) => w.abort().await,
+            FileWriter::AppendDedicated(w) => w.abort().await,
+            FileWriter::KeyValue(w) => w.abort().await,
+            FileWriter::Postpone(w) => w.abort().await,
         }
     }
 }
@@ -107,6 +119,7 @@ impl FileWriter {
 ///
 /// Call `prepare_commit()` to close all writers and collect
 /// `CommitMessage`s for use with `TableCommit`.
+/// A failed write discards pending output; create a new `TableWrite` before 
retrying.
 ///
 /// Reference: [pypaimon 
BatchTableWrite](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/write/table_write.py)
 pub struct TableWrite {
@@ -151,6 +164,8 @@ pub struct TableWrite {
     row_kind_generator: Option<RowKindGenerator>,
     row_kind_filter: Option<RowKindFilter>,
     file_index_options: Option<Arc<FileIndexOptions>>,
+    resources: Option<ResourceContext>,
+    failed: bool,
 }
 
 impl TableWrite {
@@ -425,6 +440,8 @@ impl TableWrite {
             row_kind_generator,
             row_kind_filter,
             file_index_options: file_index_options.map(Arc::new),
+            resources: None,
+            failed: false,
         })
     }
 
@@ -500,8 +517,18 @@ impl TableWrite {
         self
     }
 
+    /// Share write-buffer reservations with other consumers of this operation.
+    /// Input batches remain the caller's responsibility; charges here estimate
+    /// retained key-value batches and unflushed format-writer input. Call this
+    /// before the first write.
+    pub fn with_resources(mut self, resources: ResourceContext) -> Self {
+        self.resources = Some(resources);
+        self
+    }
+
     /// Write an Arrow RecordBatch. Rows are routed to the correct partition 
and bucket.
     pub async fn write_arrow_batch(&mut self, batch: &RecordBatch) -> 
Result<()> {
+        self.ensure_active()?;
         let Some(batch) = self.normalize_write_batch(batch)? else {
             return Ok(());
         };
@@ -850,6 +877,7 @@ impl TableWrite {
         bucket: i32,
         batch: RecordBatch,
     ) -> Result<()> {
+        self.ensure_active()?;
         let result = async {
             let key = (partition_bytes, bucket);
             if !self.partition_writers.contains_key(&key) {
@@ -862,16 +890,23 @@ impl TableWrite {
                 .await
         }
         .await;
-        if result.is_err() && self.file_index_options.is_some() {
-            for (_, writer) in self.partition_writers.drain() {
-                if let FileWriter::Append(mut writer) = writer {
-                    writer.abort().await;
-                }
-            }
+        if result.is_err() {
+            self.close().await;
+            self.failed = true;
         }
         result
     }
 
+    fn ensure_active(&self) -> Result<()> {
+        if self.failed {
+            return Err(crate::Error::DataInvalid {
+                message: "TableWrite cannot be reused after a write 
failure".to_string(),
+                source: None,
+            });
+        }
+        Ok(())
+    }
+
     /// Write multiple Arrow RecordBatches.
     pub async fn write_arrow(&mut self, batches: &[RecordBatch]) -> Result<()> 
{
         for batch in batches {
@@ -883,19 +918,15 @@ impl TableWrite {
     /// Close without preparing another commit, discarding only outstanding 
output.
     /// Files already returned by prepare_commit belong to the caller.
     pub async fn close(&mut self) {
-        for (_, writer) in self.partition_writers.drain() {
-            match writer {
-                FileWriter::Append(mut writer) => writer.abort().await,
-                FileWriter::AppendDedicated(mut writer) => 
writer.abort().await,
-                FileWriter::KeyValue(mut writer) => writer.abort().await,
-                FileWriter::Postpone(mut writer) => writer.abort().await,
-            }
+        for (_, mut writer) in self.partition_writers.drain() {
+            writer.abort().await;
         }
     }
 
     /// Close all writers and collect CommitMessages for use with TableCommit.
     /// Writers are cleared after this call, allowing the TableWrite to be 
reused.
     pub async fn prepare_commit(&mut self) -> Result<Vec<CommitMessage>> {
+        self.ensure_active()?;
         if self.file_index_options.is_some() {
             return self.prepare_indexed_append_commit().await;
         }
@@ -905,30 +936,50 @@ impl TableWrite {
         let futures: Vec<_> = writers
             .into_iter()
             .map(|((partition_bytes, bucket), writer)| async move {
-                let files = writer.prepare_commit().await?;
-                Ok::<_, crate::Error>((partition_bytes, bucket, files))
+                (partition_bytes, bucket, writer.prepare_commit().await)
             })
             .collect();
 
-        let results = futures::future::try_join_all(futures).await?;
+        let results = futures::future::join_all(futures).await;
+
+        let mut messages = Vec::new();
+        let mut error = None;
+        for (partition_bytes, bucket, result) in results {
+            match result {
+                Ok(files) if !files.data_files.is_empty() || 
!files.changelog_files.is_empty() => {
+                    let mut message = CommitMessage::new(partition_bytes, 
bucket, files.data_files);
+                    message.new_changelog_files = files.changelog_files;
+                    messages.push(message);
+                }
+                Ok(_) => {}
+                Err(err) => {
+                    error.get_or_insert(err);
+                }
+            }
+        }
+        if let Some(error) = error {
+            self.failed = true;
+            let commit = super::TableCommit::new(self.table.clone(), 
self.commit_user.clone());
+            let _ = commit.abort(&messages).await;
+            return Err(error);
+        }
 
         // Collect index files from bucket assigner
         let file_io = self.table.file_io();
-        let mut index_files_by_key = 
self.bucket_assigner.prepare_commit_index(file_io).await?;
-
-        let mut messages = Vec::new();
-        for (partition_bytes, bucket, files) in results {
-            let key = (partition_bytes.clone(), bucket);
-            let index_files = 
index_files_by_key.remove(&key).unwrap_or_default();
-            if !files.data_files.is_empty()
-                || !files.changelog_files.is_empty()
-                || !index_files.is_empty()
-            {
-                let mut msg = CommitMessage::new(partition_bytes, bucket, 
files.data_files);
-                msg.new_changelog_files = files.changelog_files;
-                msg.new_index_files = index_files;
-                messages.push(msg);
+        let mut index_files_by_key = match 
self.bucket_assigner.prepare_commit_index(file_io).await
+        {
+            Ok(files) => files,
+            Err(error) => {
+                self.failed = true;
+                let commit = super::TableCommit::new(self.table.clone(), 
self.commit_user.clone());
+                let _ = commit.abort(&messages).await;
+                return Err(error);
             }
+        };
+
+        for message in &mut messages {
+            let key = (message.partition.clone(), message.bucket);
+            message.new_index_files = 
index_files_by_key.remove(&key).unwrap_or_default();
         }
         // Emit index-only messages for (partition, bucket) pairs that had no 
data writer
         // (e.g., old buckets where keys migrated away in cross-partition 
mode).
@@ -966,6 +1017,7 @@ impl TableWrite {
             }
         }
         if let Some(error) = error {
+            self.failed = true;
             let commit = super::TableCommit::new(self.table.clone(), 
self.commit_user.clone());
             let _ = commit.abort(&messages).await;
             return Err(error);
@@ -1026,7 +1078,8 @@ impl TableWrite {
                     self.table.schema().options(),
                     &self.blob_inline_fields,
                     &self.blob_view_fields,
-                ),
+                )
+                .with_resources(self.resources.clone()),
             )))
         } else {
             Ok(FileWriter::Append(
@@ -1047,7 +1100,8 @@ impl TableWrite {
                     None,
                     None,
                 )
-                .with_file_index(self.file_index_options.clone()),
+                .with_file_index(self.file_index_options.clone())
+                .with_resources(self.resources.clone()),
             ))
         }
     }
@@ -1055,21 +1109,24 @@ impl TableWrite {
     /// Create a postpone writer (KV format, no sorting/dedup, special file 
naming).
     fn create_postpone_writer(&self, partition_path: String, bucket: i32) -> 
FileWriter {
         let data_file_prefix = format!("data-u-{}-s-0-w-", self.commit_user);
-        FileWriter::Postpone(PostponeFileWriter::new(
-            self.table.file_io().clone(),
-            PostponeWriteConfig {
-                table_location: self.table.location().to_string(),
-                partition_path,
-                bucket,
-                schema_id: self.schema_id,
-                target_file_size: self.target_file_size,
-                file_compression: self.file_compression.clone(),
-                file_compression_zstd_level: self.file_compression_zstd_level,
-                write_buffer_size: self.write_buffer_size,
-                file_format: self.file_format.clone(),
-                data_file_prefix,
-            },
-        ))
+        FileWriter::Postpone(
+            PostponeFileWriter::new(
+                self.table.file_io().clone(),
+                PostponeWriteConfig {
+                    table_location: self.table.location().to_string(),
+                    partition_path,
+                    bucket,
+                    schema_id: self.schema_id,
+                    target_file_size: self.target_file_size,
+                    file_compression: self.file_compression.clone(),
+                    file_compression_zstd_level: 
self.file_compression_zstd_level,
+                    write_buffer_size: self.write_buffer_size,
+                    file_format: self.file_format.clone(),
+                    data_file_prefix,
+                },
+            )
+            .with_resources(self.resources.clone()),
+        )
     }
 
     /// Create a key-value writer for PK tables with normal buckets.
@@ -1097,34 +1154,37 @@ impl TableWrite {
             .copied()
             .unwrap_or(0);
 
-        Ok(FileWriter::KeyValue(KeyValueFileWriter::new(
-            self.table.file_io().clone(),
-            KeyValueWriteConfig {
-                table_name: self.table.identifier().full_name(),
-                table_options: self.table.schema().options().clone(),
-                table_location: self.table.location().to_string(),
-                partition_path,
-                bucket,
-                schema_id: self.schema_id,
-                file_compression: self.file_compression.clone(),
-                file_compression_zstd_level: self.file_compression_zstd_level,
-                write_buffer_size: self.write_buffer_size,
-                file_format: self.file_format.clone(),
-                input_changelog: self.changelog_producer == 
ChangelogProducer::Input
-                    && !self.is_overwrite,
-                changelog_file_prefix: self.changelog_file_prefix.clone(),
-                changelog_file_compression: 
self.changelog_file_compression.clone(),
-                changelog_file_format: self.changelog_file_format.clone(),
-                primary_keys: self.table.schema().primary_keys().to_vec(),
-                primary_key_indices: self.primary_key_indices.clone(),
-                primary_key_types: self.primary_key_types.clone(),
-                sequence_field_indices: self.sequence_field_indices.clone(),
-                merge_engine: self.merge_engine,
-                deletion_vectors_enabled: 
CoreOptions::new(self.table.schema().options())
-                    .deletion_vectors_enabled(),
-            },
-            next_seq,
-        )?))
+        Ok(FileWriter::KeyValue(
+            KeyValueFileWriter::new(
+                self.table.file_io().clone(),
+                KeyValueWriteConfig {
+                    table_name: self.table.identifier().full_name(),
+                    table_options: self.table.schema().options().clone(),
+                    table_location: self.table.location().to_string(),
+                    partition_path,
+                    bucket,
+                    schema_id: self.schema_id,
+                    file_compression: self.file_compression.clone(),
+                    file_compression_zstd_level: 
self.file_compression_zstd_level,
+                    write_buffer_size: self.write_buffer_size,
+                    file_format: self.file_format.clone(),
+                    input_changelog: self.changelog_producer == 
ChangelogProducer::Input
+                        && !self.is_overwrite,
+                    changelog_file_prefix: self.changelog_file_prefix.clone(),
+                    changelog_file_compression: 
self.changelog_file_compression.clone(),
+                    changelog_file_format: self.changelog_file_format.clone(),
+                    primary_keys: self.table.schema().primary_keys().to_vec(),
+                    primary_key_indices: self.primary_key_indices.clone(),
+                    primary_key_types: self.primary_key_types.clone(),
+                    sequence_field_indices: 
self.sequence_field_indices.clone(),
+                    merge_engine: self.merge_engine,
+                    deletion_vectors_enabled: 
CoreOptions::new(self.table.schema().options())
+                        .deletion_vectors_enabled(),
+                },
+                next_seq,
+            )?
+            .with_resources(self.resources.clone()),
+        ))
     }
 }
 
diff --git a/crates/paimon/src/table/write_builder.rs 
b/crates/paimon/src/table/write_builder.rs
index b2af61d0..07ff11e0 100644
--- a/crates/paimon/src/table/write_builder.rs
+++ b/crates/paimon/src/table/write_builder.rs
@@ -20,6 +20,7 @@
 //! Reference: [pypaimon 
WriteBuilder](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/write/write_builder.py)
 
 use super::format_write_builder::FormatWriteBuilder;
+use crate::resource::ResourceContext;
 use crate::table::{DataEvolutionDeleteWriter, Table, TableCommit, TableUpdate, 
TableWrite};
 use uuid::Uuid;
 
@@ -75,6 +76,18 @@ impl<'a> WriteBuilder<'a> {
         }
     }
 
+    /// Share a memory budget across table writers created by `new_write`.
+    pub fn with_resources(self, resources: ResourceContext) -> Self {
+        match self.0 {
+            WriteBuilderKind::Paimon(builder) => {
+                
Self(WriteBuilderKind::Paimon(builder.with_resources(resources)))
+            }
+            WriteBuilderKind::Format(builder) => {
+                
Self(WriteBuilderKind::Format(builder.with_resources(resources)))
+            }
+        }
+    }
+
     /// Create a new TableCommit for committing write results.
     pub fn new_commit(&self) -> TableCommit {
         match &self.0 {
@@ -120,6 +133,7 @@ struct PaimonWriteBuilder<'a> {
     table: &'a Table,
     commit_user: String,
     overwrite: bool,
+    resources: Option<ResourceContext>,
 }
 
 impl<'a> PaimonWriteBuilder<'a> {
@@ -128,6 +142,7 @@ impl<'a> PaimonWriteBuilder<'a> {
             table,
             commit_user: Uuid::new_v4().to_string(),
             overwrite: false,
+            resources: None,
         }
     }
 
@@ -159,6 +174,11 @@ impl<'a> PaimonWriteBuilder<'a> {
         self
     }
 
+    pub fn with_resources(mut self, resources: ResourceContext) -> Self {
+        self.resources = Some(resources);
+        self
+    }
+
     /// Create a new TableCommit for committing write results.
     pub fn new_commit(&self) -> TableCommit {
         TableCommit::new(self.table.clone(), self.commit_user.clone())
@@ -179,7 +199,10 @@ impl<'a> PaimonWriteBuilder<'a> {
     /// when the first writer for that partition is created.
     pub fn new_write(&self) -> crate::Result<TableWrite> {
         ensure_table_write_allowed(self.table)?;
-        let write = TableWrite::new(self.table, self.commit_user.clone())?;
+        let mut write = TableWrite::new(self.table, self.commit_user.clone())?;
+        if let Some(resources) = &self.resources {
+            write = write.with_resources(resources.clone());
+        }
         Ok(if self.overwrite {
             write.with_overwrite()
         } else {
diff --git a/crates/paimon/tests/writer_resources_test.rs 
b/crates/paimon/tests/writer_resources_test.rs
new file mode 100644
index 00000000..02bd6711
--- /dev/null
+++ b/crates/paimon/tests/writer_resources_test.rs
@@ -0,0 +1,359 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod common;
+
+use common::incremental_helpers::{
+    make_batch, make_partitioned_batch, memory_table, partitioned_pk_schema, 
pk_schema, setup_dirs,
+};
+use paimon::resource::{MemoryPool, ResourceContext};
+use paimon::spec::{DataType, IntType, Schema, TableSchema};
+use paimon::Error;
+use std::sync::atomic::{AtomicBool, Ordering};
+use std::sync::Arc;
+
+#[derive(Debug)]
+struct FailOncePool(AtomicBool);
+
+impl MemoryPool for FailOncePool {
+    fn try_reserve(&self, _bytes: usize) -> paimon::Result<()> {
+        if self.0.swap(false, Ordering::SeqCst) {
+            Err(Error::ResourceExhausted {
+                message: "injected write-budget failure".to_string(),
+            })
+        } else {
+            Ok(())
+        }
+    }
+
+    fn release(&self, _bytes: usize) {}
+}
+
+async fn assert_no_data_files(io: &paimon::io::FileIO, path: &str) {
+    let files = io.list_status_recursive(path).await.unwrap();
+    assert!(
+        files.iter().all(|file| !file.path.ends_with(".parquet")),
+        "uncommitted files left behind: {files:?}"
+    );
+}
+
+#[tokio::test]
+async fn key_value_buffer_shares_limit_and_releases_on_failure() {
+    let path = "memory:/writer_resources_pk";
+    let (io, table) = memory_table(path, pk_schema(&[("file.format", 
"parquet")]));
+    setup_dirs(&io, path).await;
+    let batch = make_batch(vec![1, 2], vec![10, 20]);
+    let bytes: usize = batch
+        .columns()
+        .iter()
+        .map(|column| column.get_buffer_memory_size())
+        .sum();
+    let resources = ResourceContext::builder()
+        .memory_limit(bytes)
+        .build()
+        .unwrap();
+    let builder = table.new_write_builder().with_resources(resources.clone());
+    let mut write = builder.new_write().unwrap();
+    write.write_arrow_batch(&batch).await.unwrap();
+    assert_eq!(resources.metrics().reserved_memory_bytes, bytes);
+    assert!(matches!(
+        write.write_arrow_batch(&batch).await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert!(write.write_arrow_batch(&batch).await.is_err());
+    assert!(write.prepare_commit().await.is_err());
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_no_data_files(&io, path).await;
+}
+
+#[tokio::test]
+async fn key_value_buffer_releases_after_prepare_commit() {
+    let path = "memory:/writer_resources_pk_commit";
+    let (io, table) = memory_table(path, pk_schema(&[("file.format", 
"parquet")]));
+    setup_dirs(&io, path).await;
+    let resources = ResourceContext::builder()
+        .memory_limit(1024 * 1024)
+        .build()
+        .unwrap();
+    let builder = table.new_write_builder().with_resources(resources.clone());
+    let mut write = builder.new_write().unwrap();
+    write
+        .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    assert!(resources.metrics().reserved_memory_bytes > 0);
+    let messages = write.prepare_commit().await.unwrap();
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_eq!(messages.len(), 1);
+}
+
+#[tokio::test]
+async fn append_writer_reserves_until_commit_preparation() {
+    let path = "memory:/writer_resources_append";
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .option("file.format", "parquet")
+        .build()
+        .unwrap();
+    let (io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&io, path).await;
+    let resources = ResourceContext::builder()
+        .memory_limit(1024 * 1024)
+        .build()
+        .unwrap();
+    let builder = table.new_write_builder().with_resources(resources.clone());
+    let mut write = builder.new_write().unwrap();
+    write
+        .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    assert!(resources.metrics().reserved_memory_bytes > 0);
+    let messages = write.prepare_commit().await.unwrap();
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_eq!(messages.len(), 1);
+}
+
+#[tokio::test]
+async fn append_writer_rejects_batch_before_format_write() {
+    let path = "memory:/writer_resources_zero";
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .option("file.format", "parquet")
+        .build()
+        .unwrap();
+    let (io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&io, path).await;
+    let resources = 
ResourceContext::builder().memory_limit(0).build().unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    assert!(matches!(
+        write
+            .write_arrow_batch(&make_batch(vec![1], vec![10]))
+            .await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert!(write.prepare_commit().await.is_err());
+    assert_no_data_files(&io, path).await;
+}
+
+#[tokio::test]
+async fn postpone_writer_uses_the_same_budget() {
+    let path = "memory:/writer_resources_postpone";
+    let (io, table) = memory_table(
+        path,
+        pk_schema(&[("bucket", "-2"), ("file.format", "parquet")]),
+    );
+    setup_dirs(&io, path).await;
+    let resources = 
ResourceContext::builder().memory_limit(0).build().unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    assert!(matches!(
+        write
+            .write_arrow_batch(&make_batch(vec![1], vec![10]))
+            .await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert!(write.prepare_commit().await.is_err());
+    assert_no_data_files(&io, path).await;
+}
+
+#[tokio::test]
+async fn append_budget_rejection_prevents_partial_commit() {
+    let path = "memory:/writer_resources_append_retry";
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .option("file.format", "parquet")
+        .build()
+        .unwrap();
+    let (io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&io, path).await;
+    let limit = 1024 * 1024;
+    let resources = ResourceContext::builder()
+        .memory_limit(limit)
+        .build()
+        .unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    write
+        .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    let mut other = resources.reservation();
+    other
+        .try_grow(limit - resources.metrics().reserved_memory_bytes)
+        .unwrap();
+    assert!(matches!(
+        write
+            .write_arrow_batch(&make_batch(vec![3], vec![30]))
+            .await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    drop(other);
+    assert!(write.prepare_commit().await.is_err());
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_no_data_files(&io, path).await;
+}
+
+#[tokio::test]
+async fn postpone_budget_rejection_prevents_partial_commit() {
+    let path = "memory:/writer_resources_postpone_retry";
+    let (io, table) = memory_table(
+        path,
+        pk_schema(&[("bucket", "-2"), ("file.format", "parquet")]),
+    );
+    setup_dirs(&io, path).await;
+    let limit = 1024 * 1024;
+    let resources = ResourceContext::builder()
+        .memory_limit(limit)
+        .build()
+        .unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    write
+        .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    let mut other = resources.reservation();
+    other
+        .try_grow(limit - resources.metrics().reserved_memory_bytes)
+        .unwrap();
+    assert!(matches!(
+        write
+            .write_arrow_batch(&make_batch(vec![3], vec![30]))
+            .await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    drop(other);
+    assert!(write.prepare_commit().await.is_err());
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_no_data_files(&io, path).await;
+}
+
+#[tokio::test]
+async fn key_value_flush_rejection_prevents_partial_commit() {
+    let path = "memory:/writer_resources_pk_flush";
+    let (io, table) = memory_table(
+        path,
+        pk_schema(&[
+            ("write.parquet-buffer-size", "1b"),
+            ("file.format", "parquet"),
+        ]),
+    );
+    setup_dirs(&io, path).await;
+    let batch = make_batch(vec![1, 2], vec![10, 20]);
+    let bytes: usize = batch
+        .columns()
+        .iter()
+        .map(|column| column.get_buffer_memory_size())
+        .sum();
+    let resources = ResourceContext::builder()
+        .memory_limit(bytes)
+        .build()
+        .unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    assert!(matches!(
+        write.write_arrow_batch(&batch).await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    assert!(write.prepare_commit().await.is_err());
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_no_data_files(&io, path).await;
+}
+
+#[tokio::test]
+async fn key_value_prepare_rejection_cleans_uncommitted_file() {
+    let path = "memory:/writer_resources_pk_prepare";
+    let (io, table) = memory_table(path, pk_schema(&[("file.format", 
"parquet")]));
+    setup_dirs(&io, path).await;
+    let batch = make_batch(vec![1, 2], vec![10, 20]);
+    let bytes: usize = batch
+        .columns()
+        .iter()
+        .map(|column| column.get_buffer_memory_size())
+        .sum();
+    let resources = ResourceContext::builder()
+        .memory_limit(bytes)
+        .build()
+        .unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    write.write_arrow_batch(&batch).await.unwrap();
+    assert!(matches!(
+        write.prepare_commit().await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_no_data_files(&io, path).await;
+    assert!(write.prepare_commit().await.is_err());
+}
+
+#[tokio::test]
+async fn one_partition_prepare_failure_cleans_other_partition_outputs() {
+    let path = "memory:/writer_resources_multi_partition_prepare";
+    let (io, table) = memory_table(path, partitioned_pk_schema("1"));
+    setup_dirs(&io, path).await;
+    let pool = Arc::new(FailOncePool(AtomicBool::new(false)));
+    let resources = ResourceContext::builder()
+        .memory_pool(pool.clone())
+        .build()
+        .unwrap();
+    let mut write = table
+        .new_write_builder()
+        .with_resources(resources.clone())
+        .new_write()
+        .unwrap();
+    write
+        .write_arrow_batch(&make_partitioned_batch(
+            vec!["a", "b"],
+            vec![1, 2],
+            vec![10, 20],
+        ))
+        .await
+        .unwrap();
+    pool.0.store(true, Ordering::SeqCst);
+    assert!(matches!(
+        write.prepare_commit().await,
+        Err(Error::ResourceExhausted { .. })
+    ));
+    assert_eq!(resources.metrics().reserved_memory_bytes, 0);
+    assert_no_data_files(&io, path).await;
+}

Reply via email to