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 bb4f8bb3 feat(c): attach ResourceContext to table writers (#933)
bb4f8bb3 is described below
commit bb4f8bb362a405841cbe47b801217d31320bd1c1
Author: QuakeWang <[email protected]>
AuthorDate: Thu Sep 24 17:09:39 2026 +0800
feat(c): attach ResourceContext to table writers (#933)
---
bindings/c/src/resource.rs | 5 +-
bindings/c/src/tests.rs | 155 +++++++++++++++++++++++++++++++++++++++++++++
bindings/c/src/types.rs | 2 +
bindings/c/src/write.rs | 67 ++++++++++++++++++++
docs/src/c-binding.md | 26 ++++++++
5 files changed, 253 insertions(+), 2 deletions(-)
diff --git a/bindings/c/src/resource.rs b/bindings/c/src/resource.rs
index 4c83216a..07c9abac 100644
--- a/bindings/c/src/resource.rs
+++ b/bindings/c/src/resource.rs
@@ -23,7 +23,8 @@ use crate::error::{check_non_null, paimon_error};
use crate::result::paimon_result_resource_context;
use crate::types::{paimon_resource_context, paimon_resource_metrics};
-/// Create a shared reader reservation budget in bytes. Zero rejects nonempty
reservations.
+/// Create a shared reader and writer reservation budget in bytes.
+/// Zero rejects nonempty reservations.
#[no_mangle]
pub extern "C" fn paimon_resource_context_create(
memory_limit_bytes: usize,
@@ -73,7 +74,7 @@ pub unsafe extern "C" fn paimon_resource_context_metrics(
std::ptr::null_mut()
}
-/// Free a resource context handle. Builders and streams retain their own
clones.
+/// Free a resource context handle. Builders, streams, and writers retain
their own clones.
///
/// # Safety
/// `context` must be a handle returned by `paimon_resource_context_create`,
or null.
diff --git a/bindings/c/src/tests.rs b/bindings/c/src/tests.rs
index 872d59a0..22bba81f 100644
--- a/bindings/c/src/tests.rs
+++ b/bindings/c/src/tests.rs
@@ -1248,6 +1248,152 @@ fn
test_read_resources_share_budget_and_release_reservations() {
}
}
+#[test]
+fn test_write_builder_resources_reject_zero_budget_after_handle_free() {
+ let path = "memory:/test_write_builder_resources_zero";
+ let file_io = memory_file_io();
+ setup_table_dirs(&file_io, path);
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "test"),
+ path.to_string(),
+ simple_table_schema(),
+ None,
+ );
+ let handle = unsafe { wrap_table(table) };
+
+ unsafe {
+ let builder = paimon_table_new_write_builder(handle).write_builder;
+ let budget = paimon_resource_context_create(0).context;
+ assert!(paimon_write_builder_with_resources(builder,
budget).is_null());
+ paimon_resource_context_free(budget);
+
+ let result = paimon_write_builder_new_write(builder);
+ assert!(result.error.is_null());
+ let (array, schema) = export_batch_to_ffi(make_batch(vec![1],
vec!["a"]));
+ let error = paimon_table_write_write_arrow_batch(
+ result.write,
+ (&**array) as *const FFI_ArrowArray as *mut c_void,
+ (&**schema) as *const FFI_ArrowSchema as *mut c_void,
+ );
+ assert!(!error.is_null());
+ assert_eq!((*error).code, PaimonErrorCode::ResourceExhausted as i32);
+ paimon_error_free(error);
+
+ paimon_table_write_free(result.write);
+ paimon_write_builder_free(builder);
+ unwrap_table(handle);
+ }
+}
+
+#[test]
+fn test_write_builders_share_resource_budget() {
+ let path = "memory:/test_write_builders_share_resource_budget";
+ let file_io = memory_file_io();
+ setup_table_dirs(&file_io, path);
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "test"),
+ path.to_string(),
+ simple_table_schema(),
+ None,
+ );
+ let handle = unsafe { wrap_table(table) };
+ let value = "x".repeat(600_000);
+
+ unsafe {
+ let budget = paimon_resource_context_create(1_000_000).context;
+ let first_builder =
paimon_table_new_write_builder(handle).write_builder;
+ let second_builder =
paimon_table_new_write_builder(handle).write_builder;
+ assert!(paimon_write_builder_with_resources(first_builder,
budget).is_null());
+ assert!(paimon_write_builder_with_resources(second_builder,
budget).is_null());
+ let first = paimon_write_builder_new_write(first_builder);
+ let second = paimon_write_builder_new_write(second_builder);
+ assert!(first.error.is_null());
+ assert!(second.error.is_null());
+
+ let (array, schema) = export_batch_to_ffi(make_batch(vec![1],
vec![&value]));
+ assert!(paimon_table_write_write_arrow_batch(
+ first.write,
+ (&**array) as *const FFI_ArrowArray as *mut c_void,
+ (&**schema) as *const FFI_ArrowSchema as *mut c_void,
+ )
+ .is_null());
+ let reserved = resource_metrics(budget).reserved_memory_bytes;
+ assert!(reserved > 0 && reserved <= 1_000_000);
+
+ let (array, schema) = export_batch_to_ffi(make_batch(vec![2],
vec![&value]));
+ let error = paimon_table_write_write_arrow_batch(
+ second.write,
+ (&**array) as *const FFI_ArrowArray as *mut c_void,
+ (&**schema) as *const FFI_ArrowSchema as *mut c_void,
+ );
+ assert!(!error.is_null());
+ assert_eq!((*error).code, PaimonErrorCode::ResourceExhausted as i32);
+ paimon_error_free(error);
+ assert_eq!(resource_metrics(budget).reserved_memory_bytes, reserved);
+
+ paimon_table_write_free(first.write);
+ assert_eq!(resource_metrics(budget).reserved_memory_bytes, 0);
+ paimon_table_write_free(second.write);
+ assert_eq!(resource_metrics(budget).reserved_memory_bytes, 0);
+ assert!(resource_metrics(budget).peak_reserved_memory_bytes >=
reserved);
+ paimon_write_builder_free(first_builder);
+ paimon_write_builder_free(second_builder);
+ paimon_resource_context_free(budget);
+ unwrap_table(handle);
+ }
+}
+
+#[test]
+fn test_postpone_fixed_bucket_write_builder_resources_reject_zero_budget() {
+ let path = "memory:/test_fixed_bucket_writer_resources_zero";
+ let file_io = memory_file_io();
+ setup_table_dirs(&file_io, path);
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "test"),
+ path.to_string(),
+ partitioned_postpone_table_schema(),
+ None,
+ );
+ let handle = unsafe { wrap_table(table) };
+
+ unsafe {
+ let budget = paimon_resource_context_create(0).context;
+ let builder =
paimon_table_new_postpone_fixed_bucket_write_builder(handle).write_builder;
+ assert!(
+ paimon_postpone_fixed_bucket_write_builder_with_resources(builder,
budget).is_null()
+ );
+ let (array, schema) =
+ export_batch_to_ffi(make_postpone_bucket_plan_batch(vec!["p"],
vec![1]));
+ assert!(paimon_postpone_fixed_bucket_write_builder_with_bucket_plan(
+ builder,
+ (&**array) as *const FFI_ArrowArray as *mut c_void,
+ (&**schema) as *const FFI_ArrowSchema as *mut c_void,
+ )
+ .is_null());
+ let result =
paimon_postpone_fixed_bucket_write_builder_new_write(builder);
+ assert!(result.error.is_null());
+ let (array, schema) =
+ export_batch_to_ffi(make_partitioned_write_batch(vec!["p"],
vec![1], vec!["a"]));
+ let error = paimon_postpone_fixed_bucket_table_write_write_arrow_batch(
+ result.write,
+ (&**array) as *const FFI_ArrowArray as *mut c_void,
+ (&**schema) as *const FFI_ArrowSchema as *mut c_void,
+ );
+ assert!(!error.is_null());
+ assert_eq!((*error).code, PaimonErrorCode::ResourceExhausted as i32);
+ paimon_error_free(error);
+
+ paimon_postpone_fixed_bucket_table_write_free(result.write);
+ assert_eq!(resource_metrics(budget).reserved_memory_bytes, 0);
+ paimon_postpone_fixed_bucket_write_builder_free(builder);
+ paimon_resource_context_free(budget);
+ unwrap_table(handle);
+ }
+}
+
#[test]
fn test_read_with_projection() {
let path = "memory:/test_read_proj";
@@ -2687,6 +2833,15 @@ fn test_null_pointer_handling() {
assert!(result.write.is_null());
paimon_error_free(result.error);
+ let error = paimon_write_builder_with_resources(ptr::null_mut(),
ptr::null());
+ assert!(!error.is_null());
+ paimon_error_free(error);
+
+ let error =
+
paimon_postpone_fixed_bucket_write_builder_with_resources(ptr::null_mut(),
ptr::null());
+ assert!(!error.is_null());
+ paimon_error_free(error);
+
let result = paimon_write_builder_new_commit(ptr::null());
assert!(!result.error.is_null());
assert!(result.commit.is_null());
diff --git a/bindings/c/src/types.rs b/bindings/c/src/types.rs
index 0f95d6d0..bedd0a66 100644
--- a/bindings/c/src/types.rs
+++ b/bindings/c/src/types.rs
@@ -372,6 +372,7 @@ pub(crate) struct WriteBuilderState {
pub table: Table,
pub commit_user: String,
pub overwrite: bool,
+ pub resources: Option<ResourceContext>,
}
pub(crate) struct PostponeFixedBucketWriteBuilderState {
@@ -379,6 +380,7 @@ pub(crate) struct PostponeFixedBucketWriteBuilderState {
pub commit_user: String,
pub overwrite: bool,
pub bucket_plan: Option<PostponeBucketPlan>,
+ pub resources: Option<ResourceContext>,
}
pub(crate) struct TableWriteState {
diff --git a/bindings/c/src/write.rs b/bindings/c/src/write.rs
index c6e9afa7..86f86e18 100644
--- a/bindings/c/src/write.rs
+++ b/bindings/c/src/write.rs
@@ -22,6 +22,7 @@ use std::sync::Arc;
use arrow_array::ffi::{from_ffi, FFI_ArrowArray, FFI_ArrowSchema};
use arrow_array::{Array, RecordBatch, RecordBatchOptions, StructArray};
use arrow_schema::{DataType as ArrowDataType, Schema as ArrowSchema};
+use paimon::resource::ResourceContext;
use paimon::table::{PostponeBucketPlan, Table};
use crate::error::{check_non_null, paimon_error, validate_cstr,
PaimonErrorCode};
@@ -60,6 +61,7 @@ unsafe fn new_write_builder(
table: table_ref.clone(),
commit_user,
overwrite: false,
+ resources: None,
},
Err(error) => {
return paimon_result_write_builder {
@@ -112,6 +114,7 @@ unsafe fn new_postpone_fixed_bucket_write_builder(
commit_user,
overwrite: false,
bucket_plan: None,
+ resources: None,
};
let inner = Box::into_raw(Box::new(state)) as *mut c_void;
paimon_result_postpone_fixed_bucket_write_builder {
@@ -235,6 +238,31 @@ pub unsafe extern "C" fn
paimon_write_builder_with_overwrite(
ptr::null_mut()
}
+/// Share a resource context with writers created from this builder.
+/// The builder clones the context, so the caller may free its handle
afterward.
+/// Passing null returns an error and leaves the builder unchanged.
+///
+/// # Safety
+/// `wb` must be a valid write builder handle, or null (returns error).
+/// `context` must be a valid resource context handle, or null (returns error).
+#[no_mangle]
+pub unsafe extern "C" fn paimon_write_builder_with_resources(
+ wb: *mut paimon_write_builder,
+ context: *const paimon_resource_context,
+) -> *mut paimon_error {
+ if let Err(error) = check_non_null(wb, "wb") {
+ return error;
+ }
+ if let Err(error) = check_non_null(context, "context") {
+ return error;
+ }
+
+ let resources = &*((*context).inner as *const ResourceContext);
+ let state = &mut *((*wb).inner as *mut WriteBuilderState);
+ state.resources = Some(resources.clone());
+ ptr::null_mut()
+}
+
/// Free a postpone fixed-bucket write builder.
///
/// # Safety
@@ -270,6 +298,31 @@ pub unsafe extern "C" fn
paimon_postpone_fixed_bucket_write_builder_with_overwri
ptr::null_mut()
}
+/// Share a resource context with fixed-bucket writers created from this
builder.
+/// The builder clones the context, so the caller may free its handle
afterward.
+/// Passing null returns an error and leaves the builder unchanged.
+///
+/// # Safety
+/// `wb` must be a valid fixed-bucket builder handle, or null (returns error).
+/// `context` must be a valid resource context handle, or null (returns error).
+#[no_mangle]
+pub unsafe extern "C" fn
paimon_postpone_fixed_bucket_write_builder_with_resources(
+ wb: *mut paimon_postpone_fixed_bucket_write_builder,
+ context: *const paimon_resource_context,
+) -> *mut paimon_error {
+ if let Err(error) = check_non_null(wb, "wb") {
+ return error;
+ }
+ if let Err(error) = check_non_null(context, "context") {
+ return error;
+ }
+
+ let resources = &*((*context).inner as *const ResourceContext);
+ let state = &mut *((*wb).inner as *mut
PostponeFixedBucketWriteBuilderState);
+ state.resources = Some(resources.clone());
+ ptr::null_mut()
+}
+
/// Set a shared `partition -> total_buckets` plan.
/// The caller retains ownership when pointer or builder validation fails. Once
/// Arrow import starts, this call consumes both structs even if plan
validation
@@ -423,6 +476,9 @@ pub unsafe extern "C" fn paimon_write_builder_new_write(
if state.overwrite {
builder = builder.with_overwrite();
}
+ if let Some(resources) = &state.resources {
+ builder = builder.with_resources(resources.clone());
+ }
let result = builder.new_write().and_then(|write| {
paimon::arrow::build_target_arrow_schema(state.table.schema().fields())
.map(|schema| (Box::new(write), schema))
@@ -483,6 +539,9 @@ pub unsafe extern "C" fn
paimon_postpone_fixed_bucket_write_builder_new_write(
if state.overwrite {
builder = builder.with_overwrite();
}
+ if let Some(resources) = &state.resources {
+ builder = builder.with_resources(resources.clone());
+ }
let result = builder.new_write().and_then(|write| {
paimon::arrow::build_target_arrow_schema(state.table.schema().fields())
.map(|schema| (Box::new(write), schema))
@@ -1297,10 +1356,18 @@ const _: unsafe extern "C" fn(
paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user;
const _: unsafe extern "C" fn(*const paimon_write_builder) ->
paimon_result_table_write =
paimon_write_builder_new_write;
+const _: unsafe extern "C" fn(
+ *mut paimon_write_builder,
+ *const paimon_resource_context,
+) -> *mut paimon_error = paimon_write_builder_with_resources;
const _: unsafe extern "C" fn(
*const paimon_postpone_fixed_bucket_write_builder,
) -> paimon_result_postpone_fixed_bucket_table_write =
paimon_postpone_fixed_bucket_write_builder_new_write;
+const _: unsafe extern "C" fn(
+ *mut paimon_postpone_fixed_bucket_write_builder,
+ *const paimon_resource_context,
+) -> *mut paimon_error =
paimon_postpone_fixed_bucket_write_builder_with_resources;
const _: unsafe extern "C" fn(*const paimon_write_builder) ->
paimon_result_table_commit =
paimon_write_builder_new_commit;
const _: unsafe extern "C" fn(
diff --git a/docs/src/c-binding.md b/docs/src/c-binding.md
index 49ea5179..8cf97b4b 100644
--- a/docs/src/c-binding.md
+++ b/docs/src/c-binding.md
@@ -707,6 +707,32 @@ Writing uses a **write-then-commit** flow:
5. Pass the messages to `paimon_table_commit_commit`.
6. Free the messages, writer, committer, and builder.
+### Shared Writer Memory Reservations
+
+Attach an optional resource context to each write builder before creating its
+writer. Use `paimon_write_builder_with_resources` for a standard builder or
+`paimon_postpone_fixed_bucket_write_builder_with_resources` for a postpone
+fixed-bucket builder. The same context can be attached to multiple builders,
+including read builders, so their reservations share one limit. Both write
+builder functions clone the context; the caller may free the original C handle
+after attaching it. Each created writer retains the context it needs even if
+its builder is freed.
+
+These functions return an error for a null builder or context and leave the
+builder unchanged. A zero-byte limit rejects nonempty reservations when writing
+or preparing a commit. The existing `ResourceExhausted` error code is `6`.
+Use `paimon_resource_context_metrics` while the C handle is available to read
+the current and peak reserved bytes. The current count returns to zero after
+all writers and readers holding reservations have released them; the peak count
+remains.
+
+This is a write memory reservation interface, not a limit on total process
+memory. The current write path charges retained key-value input batches and
+unflushed format-writer input batches. Sorting, encoding, transient routing
+batches, file indexes, caller-owned Arrow input, and allocation overhead are
+not fully charged. A reservation failure does not automatically flush or retry
+a writer; handle the error as a failed write operation.
+
The input Arrow schema must match the table schema exactly, including field
count, order, names, and types. A non-nullable table field must not contain
null
values.