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;
+}