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, &current).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]]);
     }
 }
 

Reply via email to