Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 33 additions & 0 deletions bindings/c/src/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -432,6 +432,7 @@ unsafe fn new_read_builder_state(
projected_columns: None,
filter: None,
case_sensitive: true,
limit: None,
})
}

Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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()) {
Expand Down
88 changes: 88 additions & 0 deletions bindings/c/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1167,6 +1167,94 @@ 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 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(),
TableSchema::new(0, &schema),
None,
);
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>) -> 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 >= 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,
"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);
}
}

#[test]
fn test_read_resources_share_budget_and_release_reservations() {
let path = "memory:/test_read_resources";
Expand Down
2 changes: 2 additions & 0 deletions bindings/c/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -219,12 +219,14 @@ pub(crate) struct ReadBuilderState {
pub projected_columns: Option<Vec<String>>,
pub filter: Option<Predicate>,
pub case_sensitive: bool,
pub limit: Option<usize>,
}

/// Internal state for TableScan that stores table and filter.
pub(crate) struct TableScanState {
pub table: Table,
pub filter: Option<Predicate>,
pub limit: Option<usize>,
}

#[repr(C)]
Expand Down
6 changes: 6 additions & 0 deletions bindings/go/error.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
40 changes: 40 additions & 0 deletions bindings/go/read_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,26 @@ 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
// 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
}
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:
Expand Down Expand Up @@ -247,6 +267,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,
Expand Down
58 changes: 58 additions & 0 deletions bindings/go/tests/paimon_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,64 @@ 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)
}
}

// 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)

Expand Down
29 changes: 28 additions & 1 deletion crates/paimon/src/table/table_scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
Expand Down
Loading