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 f5b27f30 feat(c): expose reader memory reservations through 
ResourceContext (#931)
f5b27f30 is described below

commit f5b27f30103ee1707e8c22d595b45f986e231496
Author: QuakeWang <[email protected]>
AuthorDate: Thu Sep 24 12:17:11 2026 +0800

    feat(c): expose reader memory reservations through ResourceContext (#931)
---
 bindings/c/src/error.rs    |   2 +
 bindings/c/src/lib.rs      |   1 +
 bindings/c/src/resource.rs |  88 ++++++++++++++++++++++++++++++++++
 bindings/c/src/result.rs   |   6 +++
 bindings/c/src/table.rs    |  33 ++++++++++++-
 bindings/c/src/tests.rs    | 117 +++++++++++++++++++++++++++++++++++++++++++++
 bindings/c/src/types.rs    |  14 ++++++
 bindings/go/error.go       |  13 ++---
 docs/src/c-binding.md      |  50 +++++++++++++++++++
 9 files changed, 317 insertions(+), 7 deletions(-)

diff --git a/bindings/c/src/error.rs b/bindings/c/src/error.rs
index db82773d..e6fdf426 100644
--- a/bindings/c/src/error.rs
+++ b/bindings/c/src/error.rs
@@ -28,6 +28,7 @@ pub enum PaimonErrorCode {
     AlreadyExists = 3,
     InvalidInput = 4,
     IoError = 5,
+    ResourceExhausted = 6,
 }
 
 /// C-compatible error type.
@@ -64,6 +65,7 @@ impl paimon_error {
             | paimon::Error::DataInvalid { .. }
             | paimon::Error::IdentifierInvalid { .. } => 
PaimonErrorCode::InvalidInput,
             paimon::Error::IoUnexpected { .. } => PaimonErrorCode::IoError,
+            paimon::Error::ResourceExhausted { .. } => 
PaimonErrorCode::ResourceExhausted,
             _ => PaimonErrorCode::Unexpected,
         };
         Self::new(code, e.to_string())
diff --git a/bindings/c/src/lib.rs b/bindings/c/src/lib.rs
index 1a348454..387190dc 100644
--- a/bindings/c/src/lib.rs
+++ b/bindings/c/src/lib.rs
@@ -25,6 +25,7 @@ mod catalog;
 mod error;
 mod file_io;
 mod identifier;
+mod resource;
 mod result;
 mod table;
 #[cfg(test)]
diff --git a/bindings/c/src/resource.rs b/bindings/c/src/resource.rs
new file mode 100644
index 00000000..4c83216a
--- /dev/null
+++ b/bindings/c/src/resource.rs
@@ -0,0 +1,88 @@
+// 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.
+
+use std::ffi::c_void;
+
+use paimon::resource::ResourceContext;
+
+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.
+#[no_mangle]
+pub extern "C" fn paimon_resource_context_create(
+    memory_limit_bytes: usize,
+) -> paimon_result_resource_context {
+    match ResourceContext::builder()
+        .memory_limit(memory_limit_bytes)
+        .build()
+    {
+        Ok(resources) => paimon_result_resource_context {
+            context: Box::into_raw(Box::new(paimon_resource_context {
+                inner: Box::into_raw(Box::new(resources)) as *mut c_void,
+            })),
+            error: std::ptr::null_mut(),
+        },
+        Err(error) => paimon_result_resource_context {
+            context: std::ptr::null_mut(),
+            error: paimon_error::from_paimon(error),
+        },
+    }
+}
+
+/// Read current and peak reservation bytes from a shared context.
+///
+/// The counters are sampled independently. They are not process allocation 
metrics.
+///
+/// # Safety
+/// `context` must be a valid resource context handle, or null (returns error).
+/// `metrics` must point to writable storage, or be null (returns error).
+#[no_mangle]
+pub unsafe extern "C" fn paimon_resource_context_metrics(
+    context: *const paimon_resource_context,
+    metrics: *mut paimon_resource_metrics,
+) -> *mut paimon_error {
+    if let Err(error) = check_non_null(context, "context") {
+        return error;
+    }
+    if let Err(error) = check_non_null(metrics, "metrics") {
+        return error;
+    }
+
+    let resources = &*((*context).inner as *const ResourceContext);
+    let snapshot = resources.metrics();
+    *metrics = paimon_resource_metrics {
+        reserved_memory_bytes: snapshot.reserved_memory_bytes,
+        peak_reserved_memory_bytes: snapshot.peak_reserved_memory_bytes,
+    };
+    std::ptr::null_mut()
+}
+
+/// Free a resource context handle. Builders and streams retain their own 
clones.
+///
+/// # Safety
+/// `context` must be a handle returned by `paimon_resource_context_create`, 
or null.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_resource_context_free(context: *mut 
paimon_resource_context) {
+    if !context.is_null() {
+        let wrapper = Box::from_raw(context);
+        if !wrapper.inner.is_null() {
+            drop(Box::from_raw(wrapper.inner as *mut ResourceContext));
+        }
+    }
+}
diff --git a/bindings/c/src/result.rs b/bindings/c/src/result.rs
index 60da7c09..8572bf2d 100644
--- a/bindings/c/src/result.rs
+++ b/bindings/c/src/result.rs
@@ -96,6 +96,12 @@ pub struct paimon_result_read_builder {
     pub error: *mut paimon_error,
 }
 
+#[repr(C)]
+pub struct paimon_result_resource_context {
+    pub context: *mut paimon_resource_context,
+    pub error: *mut paimon_error,
+}
+
 #[repr(C)]
 pub struct paimon_result_table_scan {
     pub scan: *mut paimon_table_scan,
diff --git a/bindings/c/src/table.rs b/bindings/c/src/table.rs
index 26b7bbba..b0937a65 100644
--- a/bindings/c/src/table.rs
+++ b/bindings/c/src/table.rs
@@ -428,6 +428,7 @@ unsafe fn new_read_builder_state(
 
     Ok(ReadBuilderState {
         table: resolved,
+        resources: None,
         projected_columns: None,
         filter: None,
         case_sensitive: true,
@@ -516,6 +517,32 @@ pub unsafe extern "C" fn 
paimon_table_new_read_builder_with_options(
 
 // ======================= ReadBuilder ===============================
 
+/// Share a resource context with reads created from this builder.
+///
+/// The builder clones the context, so the caller may free its handle after 
this call.
+/// Passing null returns an error and leaves the builder unchanged.
+///
+/// # Safety
+/// `rb` must be a valid read 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_read_builder_with_resources(
+    rb: *mut paimon_read_builder,
+    context: *const paimon_resource_context,
+) -> *mut paimon_error {
+    if let Err(error) = check_non_null(rb, "rb") {
+        return error;
+    }
+    if let Err(error) = check_non_null(context, "context") {
+        return error;
+    }
+
+    let resources = &*((*context).inner as *const 
paimon::resource::ResourceContext);
+    let state = &mut *((*rb).inner as *mut ReadBuilderState);
+    state.resources = Some(resources.clone());
+    std::ptr::null_mut()
+}
+
 /// Free a paimon_read_builder.
 ///
 /// # Safety
@@ -710,6 +737,7 @@ pub unsafe extern "C" fn paimon_read_builder_new_read(
         Ok(table_read) => {
             let read_state = TableReadState {
                 table: state.table.clone(),
+                resources: state.resources.clone(),
                 read_type: table_read.read_type().to_vec(),
                 data_predicates: table_read.data_predicates().to_vec(),
             };
@@ -908,11 +936,14 @@ pub unsafe extern "C" fn paimon_table_read_to_arrow(
     let end = (offset.saturating_add(length)).min(all_splits.len());
     let selected = &all_splits[start..end];
 
-    let table_read = paimon::table::TableRead::new(
+    let mut table_read = paimon::table::TableRead::new(
         &state.table,
         state.read_type.clone(),
         state.data_predicates.clone(),
     );
+    if let Some(resources) = &state.resources {
+        table_read = table_read.with_resources(resources.clone());
+    }
 
     match table_read.to_arrow(selected) {
         Ok(stream) => {
diff --git a/bindings/c/src/tests.rs b/bindings/c/src/tests.rs
index a5b1acbf..872d59a0 100644
--- a/bindings/c/src/tests.rs
+++ b/bindings/c/src/tests.rs
@@ -52,6 +52,7 @@ use crate::catalog::*;
 use crate::error::*;
 use crate::file_io::*;
 use crate::identifier::*;
+use crate::resource::*;
 use crate::table::*;
 use crate::types::*;
 use crate::vector_read::*;
@@ -362,6 +363,41 @@ unsafe fn read_rows_ffi(table: *const paimon_table) -> 
Vec<(i32, String)> {
     rows
 }
 
+unsafe fn read_stream_with_resources(
+    table: *const paimon_table,
+    context: *const paimon_resource_context,
+) -> *mut paimon_record_batch_reader {
+    let builder_result = paimon_table_new_read_builder(table);
+    assert!(builder_result.error.is_null());
+    let builder = builder_result.read_builder;
+    assert!(paimon_read_builder_with_resources(builder, context).is_null());
+
+    let scan_result = paimon_read_builder_new_scan(builder);
+    assert!(scan_result.error.is_null());
+    let plan_result = paimon_table_scan_plan(scan_result.scan);
+    assert!(plan_result.error.is_null());
+    let read_result = paimon_read_builder_new_read(builder);
+    assert!(read_result.error.is_null());
+    let stream_result =
+        paimon_table_read_to_arrow(read_result.read, plan_result.plan, 0, 
usize::MAX);
+    assert!(stream_result.error.is_null());
+
+    paimon_table_read_free(read_result.read);
+    paimon_plan_free(plan_result.plan);
+    paimon_table_scan_free(scan_result.scan);
+    paimon_read_builder_free(builder);
+    stream_result.reader
+}
+
+unsafe fn resource_metrics(context: *const paimon_resource_context) -> 
paimon_resource_metrics {
+    let mut metrics = paimon_resource_metrics {
+        reserved_memory_bytes: 0,
+        peak_reserved_memory_bytes: 0,
+    };
+    assert!(paimon_resource_context_metrics(context, &mut metrics).is_null());
+    metrics
+}
+
 // =========================================================================
 //  Catalog-free table construction tests
 // =========================================================================
@@ -1131,6 +1167,87 @@ fn test_read_with_data() {
     unsafe { unwrap_table(handle) };
 }
 
+#[test]
+fn test_read_resources_share_budget_and_release_reservations() {
+    let path = "memory:/test_read_resources";
+    let file_io = memory_file_io();
+    setup_table_dirs(&file_io, path);
+    let table = Table::new(
+        file_io.clone(),
+        Identifier::new("default", "test"),
+        path.to_string(),
+        simple_table_schema(),
+        None,
+    );
+    write_data_rust(&table, &[make_batch(vec![1, 2, 3], vec!["a", "b", "c"])]);
+    let table = unsafe { wrap_table(table) };
+
+    unsafe {
+        let probe_result = paimon_resource_context_create(usize::MAX);
+        assert!(probe_result.error.is_null());
+        let probe = probe_result.context;
+        let stream = read_stream_with_resources(table, probe);
+        let first = paimon_record_batch_reader_next(stream);
+        assert!(first.error.is_null());
+        assert!(!first.batch.array.is_null());
+        let single_reader_bytes = 
resource_metrics(probe).reserved_memory_bytes;
+        assert!(single_reader_bytes > 0);
+        paimon_arrow_batch_free(first.batch);
+        paimon_record_batch_reader_free(stream);
+        assert_eq!(resource_metrics(probe).reserved_memory_bytes, 0);
+        paimon_resource_context_free(probe);
+
+        let budget_result = 
paimon_resource_context_create(single_reader_bytes);
+        assert!(budget_result.error.is_null());
+        let budget = budget_result.context;
+        let first_stream = read_stream_with_resources(table, budget);
+        let second_stream = read_stream_with_resources(table, budget);
+
+        let first = paimon_record_batch_reader_next(first_stream);
+        assert!(first.error.is_null());
+        assert!(!first.batch.array.is_null());
+        assert_eq!(
+            resource_metrics(budget).reserved_memory_bytes,
+            single_reader_bytes
+        );
+
+        let second = paimon_record_batch_reader_next(second_stream);
+        assert!(!second.error.is_null());
+        assert_eq!(
+            (*second.error).code,
+            PaimonErrorCode::ResourceExhausted as i32
+        );
+        paimon_error_free(second.error);
+        assert_eq!(
+            resource_metrics(budget).reserved_memory_bytes,
+            single_reader_bytes
+        );
+
+        paimon_arrow_batch_free(first.batch);
+        paimon_record_batch_reader_free(first_stream);
+        paimon_record_batch_reader_free(second_stream);
+        let released = resource_metrics(budget);
+        assert_eq!(released.reserved_memory_bytes, 0);
+        assert_eq!(released.peak_reserved_memory_bytes, single_reader_bytes);
+        paimon_resource_context_free(budget);
+
+        let zero_result = paimon_resource_context_create(0);
+        assert!(zero_result.error.is_null());
+        let zero_stream = read_stream_with_resources(table, 
zero_result.context);
+        paimon_resource_context_free(zero_result.context);
+        let rejected = paimon_record_batch_reader_next(zero_stream);
+        assert!(!rejected.error.is_null());
+        assert_eq!(
+            (*rejected.error).code,
+            PaimonErrorCode::ResourceExhausted as i32
+        );
+        paimon_error_free(rejected.error);
+        paimon_record_batch_reader_free(zero_stream);
+
+        unwrap_table(table);
+    }
+}
+
 #[test]
 fn test_read_with_projection() {
     let path = "memory:/test_read_proj";
diff --git a/bindings/c/src/types.rs b/bindings/c/src/types.rs
index 212a60ce..0f95d6d0 100644
--- a/bindings/c/src/types.rs
+++ b/bindings/c/src/types.rs
@@ -19,6 +19,7 @@ use std::ffi::c_void;
 use std::sync::Arc;
 
 use arrow_schema::Schema as ArrowSchema;
+use paimon::resource::ResourceContext;
 use paimon::spec::{DataField, Predicate};
 use paimon::table::{
     CommitMessage, PostponeBucketPlan, PostponeFixedBucketTableCommit,
@@ -200,9 +201,21 @@ pub struct paimon_read_builder {
     pub inner: *mut c_void,
 }
 
+#[repr(C)]
+pub struct paimon_resource_context {
+    pub inner: *mut c_void,
+}
+
+#[repr(C)]
+pub struct paimon_resource_metrics {
+    pub reserved_memory_bytes: usize,
+    pub peak_reserved_memory_bytes: usize,
+}
+
 /// Internal state for ReadBuilder that stores table, projection columns, and 
filter.
 pub(crate) struct ReadBuilderState {
     pub table: Table,
+    pub resources: Option<ResourceContext>,
     pub projected_columns: Option<Vec<String>>,
     pub filter: Option<Predicate>,
     pub case_sensitive: bool,
@@ -227,6 +240,7 @@ pub struct paimon_table_read {
 /// Internal state for TableRead that stores table, projected read type, and 
data predicates.
 pub(crate) struct TableReadState {
     pub table: Table,
+    pub resources: Option<ResourceContext>,
     pub read_type: Vec<DataField>,
     pub data_predicates: Vec<Predicate>,
 }
diff --git a/bindings/go/error.go b/bindings/go/error.go
index feb6ee9c..c85b9e08 100644
--- a/bindings/go/error.go
+++ b/bindings/go/error.go
@@ -35,12 +35,13 @@ var ErrClosed = errors.New("paimon: use of closed resource")
 type ErrorCode int32
 
 const (
-       CodeUnexpected   ErrorCode = 0
-       CodeUnsupported  ErrorCode = 1
-       CodeNotFound     ErrorCode = 2
-       CodeAlreadyExist ErrorCode = 3
-       CodeInvalidInput ErrorCode = 4
-       CodeIoError      ErrorCode = 5
+       CodeUnexpected        ErrorCode = 0
+       CodeUnsupported       ErrorCode = 1
+       CodeNotFound          ErrorCode = 2
+       CodeAlreadyExist      ErrorCode = 3
+       CodeInvalidInput      ErrorCode = 4
+       CodeIoError           ErrorCode = 5
+       CodeResourceExhausted ErrorCode = 6
 )
 
 func parseError(ctx context.Context, err *paimonError) error {
diff --git a/docs/src/c-binding.md b/docs/src/c-binding.md
index 8b216336..49ea5179 100644
--- a/docs/src/c-binding.md
+++ b/docs/src/c-binding.md
@@ -264,6 +264,56 @@ paimon_record_batch_reader_free(reader);
 paimon_table_read_free(read);
 ```
 
+### Shared Reader Memory Budget
+
+Create one resource context for reads that should share a reservation limit, 
then
+attach it to each read builder before calling `paimon_read_builder_new_read`:
+
+```c
+paimon_result_resource_context budget_result =
+    paimon_resource_context_create(64 * 1024 * 1024);
+CHECK_RESULT(budget_result);
+paimon_resource_context *budget = budget_result.context;
+
+paimon_error *error = paimon_read_builder_with_resources(read_builder, budget);
+if (error != NULL) {
+    paimon_error_free(error);
+    paimon_resource_context_free(budget);
+    goto cleanup;
+}
+
+/* Create and consume readers from read_builder here. */
+
+paimon_resource_metrics metrics;
+error = paimon_resource_context_metrics(budget, &metrics);
+if (error != NULL) {
+    paimon_error_free(error);
+    paimon_resource_context_free(budget);
+    goto cleanup;
+}
+printf("current=%zu peak=%zu\n",
+       metrics.reserved_memory_bytes,
+       metrics.peak_reserved_memory_bytes);
+
+paimon_resource_context_free(budget);
+```
+
+The builder clones the context, and each read stream retains its own clone.
+The caller may free the context handle after attaching it to the builder; keep
+the handle until after reading if metrics are needed. Reusing one context 
across
+builders makes their reservations compete for the same limit. A zero-byte limit
+rejects nonempty reservations. Admission failure is reported with
+`ResourceExhausted` (error code `6`), often from
+`paimon_record_batch_reader_next` when the stream actually reads data.
+
+`reserved_memory_bytes` is the current outstanding reservation total;
+`peak_reserved_memory_bytes` is the highest total reached by this context and
+does not reset when streams end. After all streams using a context are freed,
+the current total returns to zero. These counters track estimated reader
+working memory, including projected Parquet row groups. They are not process
+allocation or RSS measurements. Returned Arrow batches retained by the caller
+are not charged to the reader context.
+
 `paimon_table_read_to_arrow` accepts an `offset` and `length`, so separate
 workers can process disjoint contiguous ranges of the same plan. The requested
 range is clamped to the number of available splits.

Reply via email to