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 51a8ce9e feat(read): bound data evolution BLOB reads by LIMIT (#896)
51a8ce9e is described below
commit 51a8ce9e1d86743312e4a1077450781dea3ab3c7
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 16:11:25 2026 +0800
feat(read): bound data evolution BLOB reads by LIMIT (#896)
---
.../python/python/pypaimon_rust/datafusion.pyi | 5 +-
crates/paimon/src/table/data_evolution_reader.rs | 328 ++++++++++++++++++++-
crates/paimon/src/table/read_builder.rs | 10 +-
crates/paimon/src/table/table_read.rs | 16 +
crates/paimon/tests/scan_planning_parity_test.rs | 27 +-
5 files changed, 357 insertions(+), 29 deletions(-)
diff --git a/bindings/python/python/pypaimon_rust/datafusion.pyi
b/bindings/python/python/pypaimon_rust/datafusion.pyi
index dbbdaba6..4d0811ab 100644
--- a/bindings/python/python/pypaimon_rust/datafusion.pyi
+++ b/bindings/python/python/pypaimon_rust/datafusion.pyi
@@ -89,8 +89,9 @@ class ReadBuilder:
"""
...
def with_limit(self, limit: int) -> "ReadBuilder":
- """Set a scan-planning row-limit hint, not an exact cap: a matching
split is
- returned whole. Apply application-level limiting for an exact bound."""
+ """Set a scan-planning hint; data-evolution reads also stop at this
+ limit before resolving BLOB payloads. Other reads still need an
+ application-level limit for an exact bound."""
...
def with_include_row_kind(self, include: bool) -> "ReadBuilder":
"""Include a leading ``rowkind`` string column in native read
results."""
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index 8bfe68f5..5b6563f1 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -131,12 +131,22 @@ pub(crate) struct DataEvolutionReader {
blob_read_limiter: BlobReadLimiter,
blob_parallelism: usize,
batch_size: Option<usize>,
+ limit: Option<usize>,
parquet_read_budget: Option<Arc<ReadBudget>>,
table_options: Arc<HashMap<String, String>>,
mosaic_prefetch: MosaicPrefetchOptions,
read_timing: Option<Arc<DataFileReadTiming>>,
}
+fn take_limited_batch(batch: RecordBatch, remaining: &mut Option<usize>) ->
RecordBatch {
+ let Some(left) = remaining else {
+ return batch;
+ };
+ let taken = batch.num_rows().min(*left);
+ *left -= taken;
+ batch.slice(0, taken)
+}
+
impl DataEvolutionReader {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
@@ -212,6 +222,7 @@ impl DataEvolutionReader {
blob_read_limiter: BlobReadLimiter::new(),
blob_parallelism: DEFAULT_BLOB_READ_PARALLELISM,
batch_size: None,
+ limit: None,
parquet_read_budget: None,
table_options: Arc::new(HashMap::new()),
mosaic_prefetch: MosaicPrefetchOptions::default(),
@@ -231,6 +242,19 @@ impl DataEvolutionReader {
self
}
+ pub(crate) fn with_limit(mut self, limit: Option<usize>) -> Self {
+ self.limit = limit;
+ self
+ }
+
+ fn effective_batch_size(&self) -> Option<usize> {
+ match (self.batch_size, self.limit) {
+ (Some(size), Some(limit)) if limit > 0 => Some(size.min(limit)),
+ (None, Some(limit)) if limit > 0 => Some(limit),
+ (size, _) => size,
+ }
+ }
+
pub(crate) fn with_parquet_read_budget(
mut self,
parquet_read_budget: Option<Arc<ReadBudget>>,
@@ -259,9 +283,13 @@ impl DataEvolutionReader {
/// Read data files in data evolution mode.
pub fn read(self, data_splits: &[DataSplit]) ->
crate::Result<ArrowRecordBatchStream> {
+ if self.limit == Some(0) {
+ return Ok(futures::stream::empty().boxed());
+ }
let splits: Vec<DataSplit> = data_splits.to_vec();
Ok(try_stream! {
+ let mut remaining = self.limit;
let resolve_blob_views = !self.blob_view_read_fields().is_empty();
let descriptor_fields =
self.descriptor_fields_to_resolve(resolve_blob_views);
let filter_before_blob_resolution =
@@ -285,6 +313,10 @@ impl DataEvolutionReader {
let push_down_raw_predicates = !self.predicates.is_empty()
&& self.row_id_index.is_none()
&& filter_before_blob_resolution;
+ // A managed BLOB file fetches payloads as its batch is decoded.
+ // Keep batches small and restrict predicate-free file selections
+ // to the remaining output quota below.
+ let batch_size = self.effective_batch_size();
let raw_file_reader = DataFileReader::new(
self.file_io.clone(),
self.schema_manager.clone(),
@@ -297,14 +329,17 @@ impl DataEvolutionReader {
Vec::new()
},
)
- .with_batch_size(self.batch_size)
+ .with_batch_size(batch_size)
.with_blob_parallelism(self.blob_parallelism)
.with_parquet_read_budget(self.parquet_read_budget.clone())
.with_table_options(Arc::clone(&self.table_options))
.with_mosaic_prefetch(self.mosaic_prefetch)
.with_read_timing(self.read_timing.clone());
- for split in splits {
+ 'splits: for split in splits {
+ if remaining == Some(0) {
+ break;
+ }
let row_ranges = split.row_ranges().map(|r| r.to_vec());
// A chunk may span several disjoint row-id groups while one
// group still needs column-wise merging. Process each aligned
@@ -313,8 +348,14 @@ impl DataEvolutionReader {
let file_groups = reader_file_groups(split.data_files());
for files in file_groups {
+ if remaining == Some(0) {
+ break 'splits;
+ }
if is_raw_convertible(&files) {
for file_meta in files {
+ if remaining == Some(0) {
+ break 'splits;
+ }
let deletion_vector = read_file_deletion_vector(
&self.file_io,
&split,
@@ -330,7 +371,23 @@ impl DataEvolutionReader {
.await?;
let has_row_id = file_meta.first_row_id.is_some();
- let effective_row_ranges = if has_row_id {
row_ranges.clone() } else { None };
+ let mut effective_row_ranges = if has_row_id {
row_ranges.clone() } else { None };
+ if self.predicates.is_empty() {
+ if let Some(left) = remaining {
+ let selected =
selected_absolute_row_ranges_for_file(
+ file_meta.first_row_id.unwrap_or(0),
+ file_meta.row_count,
+ effective_row_ranges.as_deref(),
+ deletion_vector.as_deref(),
+ )?;
+ effective_row_ranges =
Some(prefix_row_ranges(
+ selected,
+ file_meta.first_row_id.unwrap_or(0),
+ file_meta.row_count,
+ left,
+ ));
+ }
+ }
let selected_row_ids = if
self.row_id_index.is_some() && has_row_id {
selected_absolute_row_ranges_for_file(
@@ -360,7 +417,8 @@ impl DataEvolutionReader {
deletion_vector,
effective_row_ranges,
)?;
- while let Some(batch) = stream.next().await {
+ while remaining != Some(0) {
+ let Some(batch) = stream.next().await else {
break };
let batch = batch?;
let num_rows = batch.num_rows();
let batch = if let Some(idx) =
self.row_id_index {
@@ -382,6 +440,7 @@ impl DataEvolutionReader {
blob_view_lookup.as_ref(),
&descriptor_fields,
filter_before_blob_resolution,
+ &mut remaining,
).await?;
}
}
@@ -393,15 +452,29 @@ impl DataEvolutionReader {
&prepared_group.files,
)
.await?;
- let effective_row_ranges = row_ranges.clone();
- let selected_ranges =
selected_absolute_row_ranges_for_file(
+ let mut selected_ranges =
selected_absolute_row_ranges_for_file(
prepared_group.first_row_id,
prepared_group.logical_row_count,
- effective_row_ranges.as_deref(),
+ row_ranges.as_deref(),
anchor_deletion_vector
.as_ref()
.map(|ctx| ctx.deletion_vector.as_ref()),
)?;
+ let effective_row_ranges = if
self.predicates.is_empty() {
+ if let Some(left) = remaining {
+ selected_ranges = Some(prefix_row_ranges(
+ selected_ranges,
+ prepared_group.first_row_id,
+ prepared_group.logical_row_count,
+ left,
+ ));
+ selected_ranges.clone()
+ } else {
+ row_ranges.clone()
+ }
+ } else {
+ row_ranges.clone()
+ };
let expected_output_rows = match
selected_ranges.as_ref() {
Some(ranges) => ranges.iter().map(|r| r.count() as
usize).sum(),
None => prepared_group.logical_row_count as usize,
@@ -428,7 +501,8 @@ impl DataEvolutionReader {
expected_output_rows,
anchor_deletion_vector,
)?;
- while let Some(batch) = merge_stream.next().await {
+ while remaining != Some(0) {
+ let Some(batch) = merge_stream.next().await else {
break };
let batch = batch?;
let num_rows = batch.num_rows();
let batch = if let Some(idx) = self.row_id_index {
@@ -448,6 +522,7 @@ impl DataEvolutionReader {
blob_view_lookup.as_ref(),
&descriptor_fields,
filter_before_blob_resolution,
+ &mut remaining,
).await?;
}
}
@@ -517,12 +592,16 @@ impl DataEvolutionReader {
blob_view_lookup: Option<&BlobViewLookup>,
descriptor_fields: &HashSet<String>,
filter_before_blob_resolution: bool,
+ remaining: &mut Option<usize>,
) -> crate::Result<RecordBatch> {
let mut batch = if filter_before_blob_resolution {
self.filter_wide_batch(batch)?
} else {
batch
};
+ if filter_before_blob_resolution {
+ batch = take_limited_batch(batch, remaining);
+ }
if filter_before_blob_resolution && batch.num_rows() == 0 {
return self.project_output(batch);
}
@@ -542,6 +621,7 @@ impl DataEvolutionReader {
if !filter_before_blob_resolution {
batch = self.filter_wide_batch(batch)?;
+ batch = take_limited_batch(batch, remaining);
}
self.project_output(batch)
}
@@ -673,14 +753,15 @@ impl DataEvolutionReader {
let blob_descriptor_fields = self.blob_descriptor_fields.clone();
let blob_as_descriptor = self.blob_as_descriptor;
let blob_parallelism = self.blob_parallelism;
- let batch_size = self.batch_size;
+ let batch_size = self.effective_batch_size();
let parquet_read_budget = self.parquet_read_budget.clone();
let table_options = Arc::clone(&self.table_options);
let mosaic_prefetch = self.mosaic_prefetch;
let read_timing = self.read_timing.clone();
let anchor_deletion_vector = anchor_deletion_vector.clone();
- // Batch size for column-merge output. Matches the default Parquet
reader batch size.
- const MERGE_BATCH_SIZE: usize = 1024;
+ // Match the input cap so a LIMIT cannot resolve payloads for rows that
+ // are only going to be discarded from the merge output.
+ let merge_batch_size = batch_size.unwrap_or(1024).clamp(1, 1024);
let target_schema = build_target_arrow_schema(&read_type)?;
Ok(try_stream! {
@@ -710,7 +791,7 @@ impl DataEvolutionReader {
if active_source_indices.is_empty() {
let mut emitted = 0usize;
while emitted < expected_output_rows {
- let rows_to_emit = (expected_output_rows -
emitted).min(MERGE_BATCH_SIZE);
+ let rows_to_emit = (expected_output_rows -
emitted).min(merge_batch_size);
let columns: Vec<Arc<dyn arrow_array::Array>> =
target_schema
.fields()
.iter()
@@ -840,7 +921,7 @@ impl DataEvolutionReader {
})?;
}
- let rows_to_emit = remaining.min(MERGE_BATCH_SIZE);
+ let rows_to_emit = remaining.min(merge_batch_size);
let mut columns: Vec<Arc<dyn arrow_array::Array>> =
Vec::with_capacity(source_plan.column_plan.len());
@@ -1724,6 +1805,42 @@ fn selected_absolute_row_ranges_for_file(
Ok(Some(absolute))
}
+/// Select at most the first `limit` surviving physical rows, preserving gaps
+/// from row-range selection or deletion vectors before any BLOB payload is
read.
+fn prefix_row_ranges(
+ selected: Option<Vec<RowRange>>,
+ first_row_id: i64,
+ row_count: i64,
+ limit: usize,
+) -> Vec<RowRange> {
+ let ranges = selected.unwrap_or_else(|| {
+ if row_count == 0 {
+ Vec::new()
+ } else {
+ vec![RowRange::new(first_row_id, first_row_id + row_count - 1)]
+ }
+ });
+ let mut left = limit;
+ let mut prefix = Vec::new();
+ for range in ranges {
+ let take = usize::try_from(range.count())
+ .unwrap_or(usize::MAX)
+ .min(left);
+ if take == 0 {
+ break;
+ }
+ prefix.push(RowRange::new(
+ range.from(),
+ range.from() + i64::try_from(take).unwrap_or(i64::MAX) - 1,
+ ));
+ left -= take;
+ if left == 0 {
+ break;
+ }
+ }
+ prefix
+}
+
fn non_deleted_local_ranges(row_count: i64, deletion_vector: &DeletionVector)
-> Vec<RowRange> {
let mut ranges = Vec::new();
let mut cursor = 0i64;
@@ -2742,6 +2859,24 @@ mod tests {
assert_eq!(selected, vec![RowRange::new(0, 4)]);
}
+ #[test]
+ fn test_prefix_row_ranges_counts_selected_rows_across_gaps() {
+ assert_eq!(
+ prefix_row_ranges(
+ Some(vec![RowRange::new(0, 0), RowRange::new(2, 4)]),
+ 0,
+ 5,
+ 3,
+ ),
+ vec![RowRange::new(0, 0), RowRange::new(2, 3)]
+ );
+ assert_eq!(
+ prefix_row_ranges(None, 10, 4, 3),
+ vec![RowRange::new(10, 12)]
+ );
+ assert!(prefix_row_ranges(None, 10, 4, 0).is_empty());
+ }
+
#[tokio::test]
async fn test_descriptor_columns_resolve_concurrently_and_preserve_order()
{
let schema = Arc::new(arrow_schema::Schema::new(vec![
@@ -4720,6 +4855,143 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn
test_blob_limit_skips_corrupt_payload_after_quota_across_batches_and_files() {
+ for split_blob_files in [false, true] {
+ let tempdir = tempdir().unwrap();
+ let table_path = local_file_path(tempdir.path());
+ let bucket_dir = tempdir.path().join("bucket-0");
+ fs::create_dir_all(&bucket_dir).unwrap();
+
+ let parquet_path = bucket_dir.join("data.parquet");
+ write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3,
4])], None);
+ let mut files = vec![data_file_meta_with_path(
+ "data.parquet",
+ 0,
+ 4,
+ 1,
+ parquet_path.metadata().unwrap().len() as i64,
+ Some(vec!["id"]),
+ )];
+ let blob_groups: Vec<(i64, Vec<Option<&[u8]>>)> = if
split_blob_files {
+ vec![
+ (0, vec![Some(b"first"), Some(b"second")]),
+ (2, vec![Some(b"third"), Some(b"fourth")]),
+ ]
+ } else {
+ vec![(
+ 0,
+ vec![
+ Some(b"first"),
+ Some(b"second"),
+ Some(b"third"),
+ Some(b"fourth"),
+ ],
+ )]
+ };
+ for (index, (first_row_id, values)) in
blob_groups.into_iter().enumerate() {
+ let name = format!("payload-{index}.blob");
+ let path = bucket_dir.join(&name);
+ write_blob_file(&path, &values);
+ if values.iter().any(|value| *value == Some(&b"fourth"[..])) {
+ // Preserve the index, but invalidate the fourth entry's
checksum.
+ let mut bytes = fs::read(&path).unwrap();
+ let offset = bytes
+ .windows(b"fourth".len())
+ .position(|window| window == b"fourth")
+ .map(|position| position - 4)
+ .unwrap();
+ bytes[offset] ^= 1;
+ fs::write(&path, bytes).unwrap();
+ }
+ files.push(data_file_meta_with_path(
+ &name,
+ first_row_id,
+ values.len() as i64,
+ 1,
+ path.metadata().unwrap().len() as i64,
+ Some(vec!["payload"]),
+ ));
+ }
+
+ let file_io = FileIOBuilder::new("file").build().unwrap();
+ let schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("payload", DataType::Blob(BlobType::new()))
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .option("read.batch-size", "2")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "blob_limit_batch_t"),
+ table_path,
+ schema,
+ None,
+ );
+ let raw_split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(local_file_path(&bucket_dir))
+ .with_total_buckets(1)
+ .with_data_files(files[1..].to_vec())
+ .build()
+ .unwrap();
+ let split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(local_file_path(&bucket_dir))
+ .with_total_buckets(1)
+ .with_data_files(files)
+ .build()
+ .unwrap();
+ let mut builder = table.new_read_builder();
+ builder.with_limit(3);
+ let batches = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(&[split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3]);
+ assert_eq!(
+ collect_binary_values(&batches, "payload"),
+ vec![
+ Some(b"first".to_vec()),
+ Some(b"second".to_vec()),
+ Some(b"third".to_vec()),
+ ]
+ );
+
+ // A BLOB-only raw-convertible split must obey the same quota.
+ builder.with_projection(&["payload"]).unwrap();
+ let raw_batches = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(&[raw_split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(
+ collect_binary_values(&raw_batches, "payload"),
+ vec![
+ Some(b"first".to_vec()),
+ Some(b"second".to_vec()),
+ Some(b"third".to_vec()),
+ ]
+ );
+ }
+ }
+
#[tokio::test]
async fn test_blob_fallback_defers_later_files_until_their_batch() {
use BlobFixtureValue::{Placeholder, Value};
@@ -6795,7 +7067,7 @@ mod tests {
builder.with_filter(predicate);
let read = builder.new_read().unwrap();
let batches = read
- .to_arrow(&[split])
+ .to_arrow(std::slice::from_ref(&split))
.unwrap()
.try_collect::<Vec<_>>()
.await
@@ -6803,6 +7075,17 @@ mod tests {
assert_eq!(collect_int_values(&batches, "id"), vec![2, 3, 4]);
assert_eq!(collect_int_values(&batches, "value"), vec![20, 30, 40]);
+
+ builder.with_limit(1);
+ let limited = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(&[split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&limited, "id"), vec![2]);
}
/// Multiple non-overlapping, single-file row-id segments may share one
@@ -7097,7 +7380,7 @@ mod tests {
let read = TableRead::new(&table, table.schema().fields().to_vec(),
Vec::new());
let batches = read
- .to_arrow(&[raw_split, merge_split])
+ .to_arrow(&[raw_split.clone(), merge_split.clone()])
.unwrap()
.try_collect::<Vec<_>>()
.await
@@ -7114,6 +7397,21 @@ mod tests {
collect_int_values(&batches, "id"),
vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
);
+
+ for (limit, expected) in [(0, vec![]), (1, vec![1]), (6, vec![1, 2, 3,
4, 5, 6])] {
+ let mut builder = table.new_read_builder();
+ builder.with_limit(limit);
+ let limited = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(&[raw_split.clone(), merge_split.clone()])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&limited, "id"), expected);
+ assert!(limited.iter().all(|batch| batch.num_rows() <= 2));
+ }
}
/// _ROW_ID + predicate, raw branch: surviving rows keep their ORIGINAL row
diff --git a/crates/paimon/src/table/read_builder.rs
b/crates/paimon/src/table/read_builder.rs
index 89883efa..c25f789b 100644
--- a/crates/paimon/src/table/read_builder.rs
+++ b/crates/paimon/src/table/read_builder.rs
@@ -225,7 +225,8 @@ impl<'a> ReadBuilder<'a> {
self
}
- /// Push a row-limit hint down to scan planning.
+ /// Push a row-limit hint down to scan planning. Data-evolution reads also
+ /// enforce this limit before resolving BLOB payloads.
pub fn with_limit(&mut self, limit: usize) -> &mut Self {
match &mut self.0 {
ReadBuilderKind::Paimon(builder) => {
@@ -474,9 +475,9 @@ impl<'a> PaimonReadBuilder<'a> {
/// This allows paimon-core scan planning to generate fewer splits when the
/// current scan state keeps split-level `merged_row_count()` conservative.
///
- /// Note: This method does not guarantee that exactly `limit` rows will be
- /// returned by [`TableRead`]. It is only a pushdown hint for planning.
- /// Callers or query engines are responsible for enforcing the final LIMIT.
+ /// Data-evolution [`TableRead`] enforces the limit before BLOB resolution.
+ /// Other read paths still treat it only as a planning hint, so callers or
+ /// query engines must enforce the final LIMIT themselves.
pub fn with_limit(&mut self, limit: usize) -> &mut Self {
self.limit = Some(limit);
self
@@ -543,6 +544,7 @@ impl<'a> PaimonReadBuilder<'a> {
TableRead::new(self.table, read_type,
self.filter.data_predicates.clone())
.with_parquet_read_budget(parquet_read_budget)
.with_blob_parallelism(self.blob_parallelism)
+ .map(|read| read.with_limit(self.limit))
}
/// Resolve the effective read type, deferring projection name resolution
to
diff --git a/crates/paimon/src/table/table_read.rs
b/crates/paimon/src/table/table_read.rs
index 4881fc82..ed1838d1 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -159,6 +159,19 @@ impl<'a> TableRead<'a> {
})
}
+ /// Pass the read limit to paths that can enforce it. Data-evolution reads
+ /// apply it before BLOB resolution; other Paimon reads still use the
+ /// builder limit only as a scan hint.
+ pub(crate) fn with_limit(self, limit: Option<usize>) -> Self {
+ match self.0 {
+ TableReadKind::Paimon(mut read) => {
+ read.limit = limit;
+ Self(TableReadKind::Paimon(read))
+ }
+ TableReadKind::Format(read) => Self(TableReadKind::Format(read)),
+ }
+ }
+
/// Attach an engine-specific Parquet decoder-filter factory.
///
/// The hook is used only by schema-identical raw reads. Callers must still
@@ -276,6 +289,7 @@ struct PaimonTableRead<'a> {
parquet_read_budget: Option<Arc<ReadBudget>>,
data_file_read_timing: Option<Arc<DataFileReadTiming>>,
blob_parallelism: usize,
+ limit: Option<usize>,
}
impl<'a> PaimonTableRead<'a> {
@@ -293,6 +307,7 @@ impl<'a> PaimonTableRead<'a> {
parquet_read_budget: None,
data_file_read_timing: None,
blob_parallelism: DEFAULT_BLOB_READ_PARALLELISM,
+ limit: None,
}
}
@@ -1000,6 +1015,7 @@ impl<'a> PaimonTableRead<'a> {
self.table.rest_env().cloned(),
)?
.with_batch_size(Some(core_options.read_batch_size()?))
+ .with_limit(self.limit)
.with_blob_parallelism(self.blob_parallelism)
.with_parquet_read_budget(Some(self.parquet_read_budget()?))
.with_table_options(self.table.schema().options().clone())
diff --git a/crates/paimon/tests/scan_planning_parity_test.rs
b/crates/paimon/tests/scan_planning_parity_test.rs
index 78111e39..50a408f1 100644
--- a/crates/paimon/tests/scan_planning_parity_test.rs
+++ b/crates/paimon/tests/scan_planning_parity_test.rs
@@ -352,8 +352,8 @@ async fn
position_selection_precedes_deletions_and_preserves_historical_reads()
.unwrap();
assert_eq!(plan.snapshot_id(), Some(2));
assert_eq!(read_ids(&historical, &plan).await, vec![0, 1, 2]);
- // A planning limit is a hint. It cannot drop later groups when earlier
- // selected positions are deleted; the consumer takes its final limit.
+ // The planning limit cannot drop later groups when earlier selected
+ // positions are deleted; the data-evolution reader enforces the final
limit.
let mut builder = table.new_read_builder();
builder.with_limit(2);
let plan = builder
@@ -363,8 +363,7 @@ async fn
position_selection_precedes_deletions_and_preserves_historical_reads()
.plan()
.await
.unwrap();
- let rows = read_column(&builder, &plan, 0).await;
- assert_eq!(&rows[..2], &[3, 4]);
+ assert_eq!(read_column(&builder, &plan, 0).await, vec![3, 4]);
}
#[tokio::test]
@@ -405,8 +404,13 @@ async fn
combined_delta_preserves_events_across_repeated_endpoint_deletes() {
.plan_combined_delta()
.await
.unwrap();
- // The limit is a planning hint; both selected positions belong to one
group.
- assert_eq!(read_column(&builder, &plan, 0).await, vec![4, 5]);
+ // Both selected positions belong to one planned group, but the
+ // data-evolution reader applies the requested output limit.
+ assert_eq!(
+ read_column(&table.new_read_builder(), &plan, 0).await,
+ vec![4, 5]
+ );
+ assert_eq!(read_column(&builder, &plan, 0).await, vec![4]);
let current =
table.new_read_builder().new_scan().plan().await.unwrap();
assert_eq!(read_ids(&table, ¤t).await, vec![0, 2, 3]);
}
@@ -448,7 +452,13 @@ async fn
projection_and_column_updates_do_not_multiply_positions_or_reorder_grou
.plan()
.await
.unwrap();
- assert_eq!(read_column(&builder, &plan, 0).await, vec![20, 30, 40, 999]);
+ let mut full_read = table.new_read_builder();
+ full_read.with_projection(&["value"]).unwrap();
+ assert_eq!(
+ read_column(&full_read, &plan, 0).await,
+ vec![20, 30, 40, 999]
+ );
+ assert_eq!(read_column(&builder, &plan, 0).await, vec![20]);
for (index, expected) in [vec![0, 10, 20, 30], vec![40, 999, 60, 70]]
.into_iter()
.enumerate()
@@ -460,7 +470,8 @@ async fn
projection_and_column_updates_do_not_multiply_positions_or_reorder_grou
.plan()
.await
.unwrap();
- assert_eq!(read_column(&builder, &plan, 0).await, expected);
+ assert_eq!(read_column(&full_read, &plan, 0).await, expected);
+ assert_eq!(read_column(&builder, &plan, 0).await, vec![expected[0]]);
}
}