From 58053678bcd0a68149ac59aefbf00d4a7ae4b7b2 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Thu, 24 Sep 2026 21:24:53 +0800 Subject: [PATCH 1/4] feat(go): support a read-builder limit hint for split pruning Go and C callers had no way to pass a row-count limit to scan planning, so previewing or sampling the first N rows of a large table planned every split and read every data file's statistics. The core ReadBuilder already stops selecting splits once the retained ones cover the limit (TableScan::apply_limit_pushdown); this exposes that through the bindings. Adds the paimon_read_builder_with_limit C FFI symbol and a Go ReadBuilder.WithLimit wrapper, threading the hint into both scan planning and new_read. Covered by a C-binding test (a zero limit prunes to no splits) and a Go integration test. --- bindings/c/src/table.rs | 33 ++++++++++++++++++++ bindings/c/src/tests.rs | 52 ++++++++++++++++++++++++++++++++ bindings/c/src/types.rs | 2 ++ bindings/go/read_builder.go | 32 ++++++++++++++++++++ bindings/go/tests/paimon_test.go | 42 ++++++++++++++++++++++++++ 5 files changed, 161 insertions(+) diff --git a/bindings/c/src/table.rs b/bindings/c/src/table.rs index b0937a65e..0feba23be 100644 --- a/bindings/c/src/table.rs +++ b/bindings/c/src/table.rs @@ -432,6 +432,7 @@ unsafe fn new_read_builder_state( projected_columns: None, filter: None, case_sensitive: true, + limit: None, }) } @@ -643,6 +644,29 @@ pub unsafe extern "C" fn paimon_read_builder_with_case_sensitive( std::ptr::null_mut() } +/// Set a row-count limit hint for scan planning. +/// +/// This lets planning stop selecting splits once the retained splits already +/// cover `limit` rows, cutting the plan-time I/O of listing every split's +/// stats when only the first N rows are wanted. It is a hint only: it does not +/// guarantee exactly `limit` rows are returned, so the caller still enforces +/// the final row limit when reading. +/// +/// # Safety +/// `rb` must be a valid pointer from `paimon_table_new_read_builder`, or null (returns error). +#[no_mangle] +pub unsafe extern "C" fn paimon_read_builder_with_limit( + rb: *mut paimon_read_builder, + limit: usize, +) -> *mut paimon_error { + if let Err(e) = check_non_null(rb, "rb") { + return e; + } + let state = &mut *((*rb).inner as *mut ReadBuilderState); + state.limit = Some(limit); + std::ptr::null_mut() +} + /// Set a filter predicate for scan planning. /// /// The predicate is consumed (ownership transferred to the read builder). @@ -691,6 +715,7 @@ pub unsafe extern "C" fn paimon_read_builder_new_scan( let scan_state = TableScanState { table: state.table.clone(), filter: state.filter.clone(), + limit: state.limit, }; let inner = Box::into_raw(Box::new(scan_state)) as *mut c_void; paimon_result_table_scan { @@ -733,6 +758,11 @@ pub unsafe extern "C" fn paimon_read_builder_new_read( rb_rust.with_filter(filter.clone()); } + // Apply limit hint if set + if let Some(limit) = state.limit { + rb_rust.with_limit(limit); + } + match rb_rust.new_read() { Ok(table_read) => { let read_state = TableReadState { @@ -788,6 +818,9 @@ pub unsafe extern "C" fn paimon_table_scan_plan( if let Some(ref filter) = scan_state.filter { rb.with_filter(filter.clone()); } + if let Some(limit) = scan_state.limit { + rb.with_limit(limit); + } let table_scan = rb.new_scan(); match runtime().block_on(table_scan.plan()) { diff --git a/bindings/c/src/tests.rs b/bindings/c/src/tests.rs index 22bba81f8..e13e75bb2 100644 --- a/bindings/c/src/tests.rs +++ b/bindings/c/src/tests.rs @@ -1167,6 +1167,58 @@ fn test_read_with_data() { unsafe { unwrap_table(handle) }; } +#[test] +fn test_read_builder_with_limit_prunes_plan_splits() { + // `with_limit` is a plan-time hint: planning stops selecting splits once the + // retained ones already cover the limit (mirrors core `apply_limit_pushdown`). + // Assert the FFI threads the hint through to planning -- a zero limit prunes + // every split, a generous limit leaves the plan untouched. Without the hint + // reaching planning, the zero-limit plan would keep the baseline splits. + let path = "memory:/test_read_limit"; + 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 handle = unsafe { wrap_table(table) }; + + unsafe fn plan_split_count(handle: *const paimon_table, limit: Option) -> usize { + let rb = paimon_table_new_read_builder(handle).read_builder; + if let Some(limit) = limit { + assert!(paimon_read_builder_with_limit(rb, limit).is_null()); + } + let scan = paimon_read_builder_new_scan(rb).scan; + let plan_result = paimon_table_scan_plan(scan); + assert!(plan_result.error.is_null()); + let count = paimon_plan_num_splits(plan_result.plan); + paimon_plan_free(plan_result.plan); + paimon_table_scan_free(scan); + paimon_read_builder_free(rb); + count + } + + unsafe { + let baseline = plan_split_count(handle, None); + assert!(baseline >= 1, "a table with rows should plan >= 1 split"); + assert_eq!( + plan_split_count(handle, Some(0)), + 0, + "a zero limit must prune every split at plan time" + ); + assert_eq!( + plan_split_count(handle, Some(1000)), + baseline, + "a limit above the row count must not prune any split" + ); + unwrap_table(handle); + } +} + #[test] fn test_read_resources_share_budget_and_release_reservations() { let path = "memory:/test_read_resources"; diff --git a/bindings/c/src/types.rs b/bindings/c/src/types.rs index bedd0a664..3f96b13c3 100644 --- a/bindings/c/src/types.rs +++ b/bindings/c/src/types.rs @@ -219,12 +219,14 @@ pub(crate) struct ReadBuilderState { pub projected_columns: Option>, pub filter: Option, pub case_sensitive: bool, + pub limit: Option, } /// Internal state for TableScan that stores table and filter. pub(crate) struct TableScanState { pub table: Table, pub filter: Option, + pub limit: Option, } #[repr(C)] diff --git a/bindings/go/read_builder.go b/bindings/go/read_builder.go index 3f52c3c71..5c5437877 100644 --- a/bindings/go/read_builder.go +++ b/bindings/go/read_builder.go @@ -92,6 +92,18 @@ func (rb *ReadBuilder) WithCaseSensitive(caseSensitive bool) error { return ffiReadBuilderWithCaseSensitive.symbol(rb.ctx)(rb.inner, caseSensitive) } +// WithLimit sets a row-count limit hint for scan planning. Planning stops +// selecting splits once the retained ones already cover the limit, so a preview +// or sample of the first N rows avoids listing every split's statistics. It is a +// hint only: it does not guarantee exactly limit rows are returned, so callers +// still enforce the final row limit when reading. +func (rb *ReadBuilder) WithLimit(limit int) error { + if rb.inner == nil { + return ErrClosed + } + return ffiReadBuilderWithLimit.symbol(rb.ctx)(rb.inner, uintptr(limit)) +} + // WithFilter sets a filter predicate for scan planning and read-side pruning. // // The predicate is used in two phases: @@ -247,6 +259,26 @@ var ffiReadBuilderWithCaseSensitive = newFFI(ffiOpts{ } }) +// size_t is passed pointer-width, matching ffiTableReadToArrow's offset/length. +var ffiReadBuilderWithLimit = newFFI(ffiOpts{ + sym: "paimon_read_builder_with_limit", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(rb *paimonReadBuilder, limit uintptr) error { + return func(rb *paimonReadBuilder, limit uintptr) error { + var errPtr *paimonError + ffiCall( + unsafe.Pointer(&errPtr), + unsafe.Pointer(&rb), + unsafe.Pointer(&limit), + ) + if errPtr != nil { + return parseError(ctx, errPtr) + } + return nil + } +}) + var ffiReadBuilderWithFilter = newFFI(ffiOpts{ sym: "paimon_read_builder_with_filter", rType: &ffi.TypePointer, diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index 169680940..97d339ab0 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -287,6 +287,48 @@ func openTestTable(t *testing.T) *paimon.Table { return openTableAt(t, warehouse, "simple_log_table") } +func TestReadBuilderWithLimitPrunesPlanSplits(t *testing.T) { + table := openTestTable(t) + + planSplitCount := func(setLimit func(*paimon.ReadBuilder) error) int { + rb, err := table.NewReadBuilder() + if err != nil { + t.Fatalf("Failed to create read builder: %v", err) + } + defer rb.Close() + if setLimit != nil { + if err := setLimit(rb); err != nil { + t.Fatalf("Failed to set limit: %v", err) + } + } + scan, err := rb.NewScan() + if err != nil { + t.Fatalf("Failed to create scan: %v", err) + } + defer scan.Close() + plan, err := scan.Plan() + if err != nil { + t.Fatalf("Failed to plan: %v", err) + } + defer plan.Close() + return len(plan.Splits()) + } + + baseline := planSplitCount(nil) + if baseline < 1 { + t.Fatalf("expected >= 1 split without a limit, got %d", baseline) + } + // A zero limit prunes every split at plan time, proving the hint reaches + // planning through the FFI (mirrors core apply_limit_pushdown semantics). + if got := planSplitCount(func(rb *paimon.ReadBuilder) error { return rb.WithLimit(0) }); got != 0 { + t.Fatalf("expected 0 splits with WithLimit(0), got %d", got) + } + // A limit far above the row count must not prune any split. + if got := planSplitCount(func(rb *paimon.ReadBuilder) error { return rb.WithLimit(1 << 30) }); got != baseline { + t.Fatalf("expected %d splits with a large limit, got %d", baseline, got) + } +} + func TestWriteCommitReadRoundTrip(t *testing.T) { table := openCopiedTestTable(t) From 696a5d4f179aa5acea059ff290f34c8648f7a59d Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Sat, 26 Sep 2026 09:58:08 +0800 Subject: [PATCH 2/4] fix(go): reject a negative read-builder limit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-up (#946): `ReadBuilder.WithLimit(-1)` returned success and passed `uintptr(-1)` across the C boundary, so `LimitPushdownAccumulator` saw `limit as i64 == -1` and stopped planning after the first split with a known row count — a silently truncated scan. Reject a negative limit up front with `ErrNegativeLimit`. Also correct the doc claim: planning still reads the manifest entries, so the hint bounds retained splits rather than avoiding all planning I/O. --- bindings/go/error.go | 6 ++++++ bindings/go/read_builder.go | 16 ++++++++++++---- bindings/go/tests/paimon_test.go | 16 ++++++++++++++++ 3 files changed, 34 insertions(+), 4 deletions(-) diff --git a/bindings/go/error.go b/bindings/go/error.go index c85b9e08f..af0058419 100644 --- a/bindings/go/error.go +++ b/bindings/go/error.go @@ -31,6 +31,12 @@ import ( // ErrClosed is returned when an operation is attempted on a closed resource. var ErrClosed = errors.New("paimon: use of closed resource") +// ErrNegativeLimit is returned by ReadBuilder.WithLimit when the limit is +// negative. A negative Go int would wrap to a huge value through the unsigned +// C boundary and make scan planning stop after the first split, silently +// dropping rows, so it is rejected up front. +var ErrNegativeLimit = errors.New("paimon: read-builder limit must not be negative") + // ErrorCode represents categories of errors from paimon. type ErrorCode int32 diff --git a/bindings/go/read_builder.go b/bindings/go/read_builder.go index 5c5437877..8a2e36530 100644 --- a/bindings/go/read_builder.go +++ b/bindings/go/read_builder.go @@ -93,11 +93,19 @@ func (rb *ReadBuilder) WithCaseSensitive(caseSensitive bool) error { } // WithLimit sets a row-count limit hint for scan planning. Planning stops -// selecting splits once the retained ones already cover the limit, so a preview -// or sample of the first N rows avoids listing every split's statistics. It is a -// hint only: it does not guarantee exactly limit rows are returned, so callers -// still enforce the final row limit when reading. +// retaining splits once the ones already kept cover the limit, so reading a +// preview or sample of the first N rows keeps fewer splits. Planning still reads +// the manifest entries to learn each split's row count, so this bounds how many +// splits are retained rather than avoiding all planning I/O. It is a hint only: +// it does not guarantee exactly limit rows are returned, so callers still enforce +// the final row limit when reading. +// +// A negative limit is rejected with ErrNegativeLimit: it would wrap to a huge +// value through the unsigned C boundary and stop planning after the first split. func (rb *ReadBuilder) WithLimit(limit int) error { + if limit < 0 { + return ErrNegativeLimit + } if rb.inner == nil { return ErrClosed } diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index 97d339ab0..ebb5a5efe 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -329,6 +329,22 @@ func TestReadBuilderWithLimitPrunesPlanSplits(t *testing.T) { } } +// A negative limit must be rejected before it reaches the unsigned C boundary, +// where it would wrap to a huge value and make planning stop after the first +// split. The check runs ahead of the closed-builder guard, so it needs no +// warehouse: the distinct ErrNegativeLimit (not ErrClosed) proves the reject +// arm ran rather than the nil-inner arm. +func TestReadBuilderWithLimitRejectsNegative(t *testing.T) { + rb := &paimon.ReadBuilder{} + err := rb.WithLimit(-1) + if !errors.Is(err, paimon.ErrNegativeLimit) { + t.Fatalf("expected ErrNegativeLimit for a negative limit, got %v", err) + } + if errors.Is(err, paimon.ErrClosed) { + t.Fatalf("negative limit must be rejected as invalid input, not as a closed builder") + } +} + func TestWriteCommitReadRoundTrip(t *testing.T) { table := openCopiedTestTable(t) From 4afe68e8f2da1d07ee8f8b1c07a1c633baeb641d Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Wed, 30 Sep 2026 21:12:45 +0800 Subject: [PATCH 3/4] test(c): cover the read-builder limit pruning multiple known-count splits Extend the C-FFI limit test requested in review to plan multiple splits: a tiny `source.split.target-size` keeps each committed file as its own split, so a small positive limit is shown to retain a strict, non-empty subset rather than every split or none. --- bindings/c/src/tests.rs | 39 +++++++++++++++++++++++++++++++-------- 1 file changed, 31 insertions(+), 8 deletions(-) diff --git a/bindings/c/src/tests.rs b/bindings/c/src/tests.rs index e13e75bb2..39c9794ce 100644 --- a/bindings/c/src/tests.rs +++ b/bindings/c/src/tests.rs @@ -1169,22 +1169,33 @@ fn test_read_with_data() { #[test] fn test_read_builder_with_limit_prunes_plan_splits() { - // `with_limit` is a plan-time hint: planning stops selecting splits once the - // retained ones already cover the limit (mirrors core `apply_limit_pushdown`). - // Assert the FFI threads the hint through to planning -- a zero limit prunes - // every split, a generous limit leaves the plan untouched. Without the hint - // reaching planning, the zero-limit plan would keep the baseline splits. + // `with_limit` is a plan-time hint: planning stops retaining splits once the + // kept ones already cover the limit (mirrors core `apply_limit_pushdown`). + // Three separate commits give multiple splits with known row counts, so this + // asserts the FFI threads the hint through to planning across splits: a zero + // limit prunes every split, a small limit keeps a strict non-empty subset, + // and a generous limit leaves the plan untouched. let path = "memory:/test_read_limit"; let file_io = memory_file_io(); setup_table_dirs(&file_io, path); + // Force a tiny split target so each committed file becomes its own split, + // giving multiple known-count splits instead of one bundled split. + let schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("name", DataType::VarChar(VarCharType::string_type())) + .option("source.split.target-size", "1") + .build() + .unwrap(); let table = Table::new( file_io.clone(), Identifier::new("default", "test"), path.to_string(), - simple_table_schema(), + TableSchema::new(0, &schema), None, ); - write_data_rust(&table, &[make_batch(vec![1, 2, 3], vec!["a", "b", "c"])]); + write_data_rust(&table, &[make_batch(vec![1, 2], vec!["a", "b"])]); + write_data_rust(&table, &[make_batch(vec![3, 4], vec!["c", "d"])]); + write_data_rust(&table, &[make_batch(vec![5, 6], vec!["e", "f"])]); let handle = unsafe { wrap_table(table) }; unsafe fn plan_split_count(handle: *const paimon_table, limit: Option) -> usize { @@ -1204,12 +1215,24 @@ fn test_read_builder_with_limit_prunes_plan_splits() { unsafe { let baseline = plan_split_count(handle, None); - assert!(baseline >= 1, "a table with rows should plan >= 1 split"); + assert!( + baseline >= 2, + "three commits should plan multiple splits, got {baseline}" + ); assert_eq!( plan_split_count(handle, Some(0)), 0, "a zero limit must prune every split at plan time" ); + // A small positive limit is covered by the first retained split alone, so + // planning keeps a strict non-empty subset instead of every split. A limit + // that wrapped to a huge value (the negative-input bug) would instead keep + // the full baseline here. + let pruned = plan_split_count(handle, Some(1)); + assert!( + pruned >= 1 && pruned < baseline, + "limit 1 should keep a strict non-empty subset of {baseline} splits, got {pruned}" + ); assert_eq!( plan_split_count(handle, Some(1000)), baseline, From af65feb357d6e01aabeaf92e85e27cc2627c8c74 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Fri, 2 Oct 2026 09:23:17 +0800 Subject: [PATCH 4/4] fix(core): compare the limit-pushdown accumulator without a lossy signed cast MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `LimitPushdownAccumulator::push` compared `scanned_row_count >= self.limit as i64`. `limit` is a `usize`, so any value above `i64::MAX` — e.g. a C caller passing `SIZE_MAX` through `paimon_read_builder_with_limit` — turned negative, and the first counted split satisfied the limit, silently dropping the rest. Compare as `u128` instead; `scanned_row_count` is always non-negative, so a limit a real row count cannot reach never early-stops the scan. Covered by a core accumulator test for `usize::MAX` and `i64::MAX + 1`, and extended the C FFI plan test to assert both oversized limits retain every split. --- bindings/c/src/tests.rs | 13 ++++++++++++ crates/paimon/src/table/table_scan.rs | 29 ++++++++++++++++++++++++++- 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/bindings/c/src/tests.rs b/bindings/c/src/tests.rs index 39c9794ce..1e71a908d 100644 --- a/bindings/c/src/tests.rs +++ b/bindings/c/src/tests.rs @@ -1238,6 +1238,19 @@ fn test_read_builder_with_limit_prunes_plan_splits() { baseline, "a limit above the row count must not prune any split" ); + // A limit above i64::MAX must not wrap negative in the accumulator's + // signed comparison and truncate the scan: SIZE_MAX and i64::MAX + 1 + // keep every split, like any limit the row count cannot reach. + assert_eq!( + plan_split_count(handle, Some(usize::MAX)), + baseline, + "SIZE_MAX must not truncate the plan" + ); + assert_eq!( + plan_split_count(handle, Some((i64::MAX as usize) + 1)), + baseline, + "i64::MAX + 1 must not truncate the plan" + ); unwrap_table(handle); } } diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index ceb4bc856..a3972cd4d 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -782,7 +782,12 @@ impl LimitPushdownAccumulator { self.fallback_splits.push(split.clone()); self.limited_splits.push(split); self.scanned_row_count += merged_count; - self.limit_early_stopped = self.scanned_row_count >= self.limit as i64; + // Compare without a lossy signed cast: `limit` is `usize`, so + // `limit as i64` turns any value above `i64::MAX` (e.g. a C caller + // passing `SIZE_MAX`) negative, which would early-stop after the + // first counted split. `scanned_row_count` is always >= 0, so an + // unsigned comparison keeps a huge limit from truncating the scan. + self.limit_early_stopped = self.scanned_row_count as u128 >= self.limit as u128; } else { self.fallback_splits.push(split); } @@ -3522,6 +3527,28 @@ mod tests { ); } + #[test] + fn test_incremental_limit_accumulator_does_not_truncate_on_oversized_limit() { + // A limit above i64::MAX (e.g. a C caller passing SIZE_MAX) must not wrap + // negative through a signed cast and early-stop after the first counted + // split; real row counts can never reach it, so every split is kept. + for limit in [usize::MAX, (i64::MAX as usize) + 1] { + let mut accumulator = LimitPushdownAccumulator::new(limit); + assert!(!accumulator.push(limit_test_split("a.parquet", 2))); + assert!(!accumulator.push(limit_test_split("b.parquet", 3))); + let result = accumulator.finish(); + assert!( + !result.limit_early_stopped, + "limit {limit} must not early-stop" + ); + assert_eq!( + split_file_names(&result.splits), + vec!["a.parquet", "b.parquet"], + "limit {limit} must keep all splits" + ); + } + } + #[test] fn test_incremental_limit_accumulator_returns_fallback_when_limit_not_reached() { let mut accumulator = LimitPushdownAccumulator::new(100);