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.