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 9122eb19 feat(read): avoid unselected managed BLOB payloads under 
LIMIT (#901)
9122eb19 is described below

commit 9122eb19b5d352c7c418abc5664de2de8076a232
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 18:28:28 2026 +0800

    feat(read): avoid unselected managed BLOB payloads under LIMIT (#901)
---
 crates/paimon/src/table/data_evolution_reader.rs | 192 ++++++++++++++++++++++-
 crates/paimon/src/table/source.rs                |   8 +
 2 files changed, 199 insertions(+), 1 deletion(-)

diff --git a/crates/paimon/src/table/data_evolution_reader.rs 
b/crates/paimon/src/table/data_evolution_reader.rs
index 411ba255..52de22c5 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -255,6 +255,118 @@ impl DataEvolutionReader {
         }
     }
 
+    /// Managed BLOB payloads are fetched by the format reader, before the
+    /// residual predicate. A LIMIT with a predicate on ordinary columns must
+    /// therefore select matching row IDs before opening those payloads.
+    fn needs_managed_blob_preselection(&self, splits: &[DataSplit]) -> bool {
+        if self.limit.is_none()
+            || self.predicates.is_empty()
+            || self.blob_as_descriptor
+            || !self.blob_view_read_fields().is_empty()
+            || !self
+                .wide_file_read_type
+                .iter()
+                .any(|field| field.data_type().is_blob_file_field())
+            || !splits.iter().any(|split| {
+                split
+                    .data_files()
+                    .iter()
+                    .any(|file| is_blob_file_name(&file.file_name))
+            })
+            // A legacy raw file without first_row_id cannot be addressed by
+            // the preselection's global row ranges. Keep its existing read
+            // behavior rather than changing a valid read into an error.
+            || !splits.iter().all(|split| {
+                split
+                    .data_files()
+                    .iter()
+                    .all(|file| file.first_row_id.is_some())
+            })
+        {
+            return false;
+        }
+
+        let blob_fields = self
+            .table_fields
+            .iter()
+            .filter(|field| field.data_type().is_blob_file_field())
+            .map(|field| field.name().to_string())
+            .collect::<HashSet<_>>();
+        !predicates_reference_any_field(&self.predicates, &blob_fields, 
&self.table_fields)
+    }
+
+    async fn preselect_managed_blob_rows(
+        &self,
+        splits: &[DataSplit],
+    ) -> crate::Result<Vec<DataSplit>> {
+        let Some(mut remaining) = self.limit else {
+            return Ok(splits.to_vec());
+        };
+        let mut selected = Vec::new();
+        for split in splits {
+            if remaining == 0 {
+                break;
+            }
+            // This read projects only synthesized row IDs. Predicate columns
+            // are widened internally; no managed BLOB field is requested, so
+            // a .blob source may read its index but never its payload.
+            let preselection = Self::new(
+                self.file_io.clone(),
+                self.schema_manager.clone(),
+                self.table_schema_id,
+                self.table_fields.clone(),
+                vec![crate::spec::row_id_data_field()],
+                self.predicates.clone(),
+                true,
+                HashSet::new(),
+                HashSet::new(),
+                false,
+                None,
+            )?
+            .with_batch_size(self.batch_size)
+            .with_limit(Some(remaining))
+            .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());
+            let mut rows = preselection.read(std::slice::from_ref(split))?;
+            let mut ranges = Vec::new();
+            let mut selected_count = 0usize;
+            while let Some(batch) = rows.next().await {
+                let batch = batch?;
+                let row_ids = batch
+                    .column(0)
+                    .as_any()
+                    .downcast_ref::<Int64Array>()
+                    .ok_or_else(|| Error::DataInvalid {
+                        message: "Managed BLOB preselection did not return 
_ROW_ID".to_string(),
+                        source: None,
+                    })?;
+                for index in 0..row_ids.len() {
+                    if row_ids.is_null(index) {
+                        return Err(Error::DataInvalid {
+                            message: "Managed BLOB preselection requires 
non-null _ROW_ID"
+                                .to_string(),
+                            source: None,
+                        });
+                    }
+                    let id = row_ids.value(index);
+                    ranges.push(RowRange::new(id, id));
+                    selected_count += 1;
+                }
+            }
+            if !ranges.is_empty() {
+                selected.push(
+                    split
+                        .clone()
+                        
.with_selected_row_ranges(crate::table::merge_row_ranges(ranges)),
+                );
+            }
+            remaining = remaining.saturating_sub(selected_count);
+        }
+        Ok(selected)
+    }
+
     pub(crate) fn with_parquet_read_budget(
         mut self,
         parquet_read_budget: Option<Arc<ReadBudget>>,
@@ -288,6 +400,19 @@ impl DataEvolutionReader {
         }
         let splits: Vec<DataSplit> = data_splits.to_vec();
 
+        if self.needs_managed_blob_preselection(&splits) {
+            return Ok(try_stream! {
+                let selected = 
self.preselect_managed_blob_rows(&splits).await?;
+                let mut reader = self;
+                reader.predicates.clear();
+                let mut rows = reader.read(&selected)?;
+                while let Some(batch) = rows.next().await {
+                    yield batch?;
+                }
+            }
+            .boxed());
+        }
+
         Ok(try_stream! {
             let mut remaining = self.limit;
             let resolve_blob_views = !self.blob_view_read_fields().is_empty();
@@ -4987,7 +5112,7 @@ mod tests {
             let batches = builder
                 .new_read()
                 .unwrap()
-                .to_arrow(&[split])
+                .to_arrow(std::slice::from_ref(&split))
                 .unwrap()
                 .try_collect::<Vec<_>>()
                 .await
@@ -5002,6 +5127,71 @@ mod tests {
                 ]
             );
 
+            // The first row fails an ordinary-column predicate. The second
+            // output batch must read row 3, but not the corrupt row 4 payload.
+            let mut filtered_builder = table.new_read_builder();
+            filtered_builder.with_limit(2);
+            filtered_builder.with_filter(
+                PredicateBuilder::new(table.schema().fields())
+                    .greater_than("id", Datum::Int(1))
+                    .unwrap(),
+            );
+            let filtered = filtered_builder
+                .new_read()
+                .unwrap()
+                .to_arrow(std::slice::from_ref(&split))
+                .unwrap()
+                .try_collect::<Vec<_>>()
+                .await
+                .unwrap();
+            assert_eq!(collect_int_values(&filtered, "id"), vec![2, 3]);
+            assert_eq!(
+                collect_binary_values(&filtered, "payload"),
+                vec![Some(b"second".to_vec()), Some(b"third".to_vec())]
+            );
+
+            // The same selection must preserve the quota across split
+            // boundaries, without touching row 4 in the second split.
+            let first_half = split
+                .clone()
+                .with_selected_row_ranges(vec![RowRange::new(0, 1)]);
+            let second_half = split
+                .clone()
+                .with_selected_row_ranges(vec![RowRange::new(2, 3)]);
+            let filtered_across_splits = filtered_builder
+                .new_read()
+                .unwrap()
+                .to_arrow(&[first_half, second_half])
+                .unwrap()
+                .try_collect::<Vec<_>>()
+                .await
+                .unwrap();
+            assert_eq!(
+                collect_int_values(&filtered_across_splits, "id"),
+                vec![2, 3]
+            );
+            assert_eq!(
+                collect_binary_values(&filtered_across_splits, "payload"),
+                vec![Some(b"second".to_vec()), Some(b"third".to_vec())]
+            );
+
+            let mut no_match_builder = table.new_read_builder();
+            no_match_builder.with_limit(1);
+            no_match_builder.with_filter(
+                PredicateBuilder::new(table.schema().fields())
+                    .greater_than("id", Datum::Int(4))
+                    .unwrap(),
+            );
+            let no_match = no_match_builder
+                .new_read()
+                .unwrap()
+                .to_arrow(&[split])
+                .unwrap()
+                .try_collect::<Vec<_>>()
+                .await
+                .unwrap();
+            assert!(collect_int_values(&no_match, "id").is_empty());
+
             // A BLOB-only raw-convertible split must obey the same quota.
             builder.with_projection(&["payload"]).unwrap();
             let raw_batches = builder
diff --git a/crates/paimon/src/table/source.rs 
b/crates/paimon/src/table/source.rs
index e4c30866..f96a4f47 100644
--- a/crates/paimon/src/table/source.rs
+++ b/crates/paimon/src/table/source.rs
@@ -544,6 +544,14 @@ impl DataSplit {
         self.row_ranges.as_deref()
     }
 
+    /// Narrow an already planned split without rebuilding or copying its file
+    /// and deletion-vector metadata. Ranges use the same coordinate system as
+    /// the existing `row_ranges` field.
+    pub(crate) fn with_selected_row_ranges(mut self, ranges: Vec<RowRange>) -> 
Self {
+        self.row_ranges = Some(ranges.into());
+        self
+    }
+
     /// Whether this split can be read raw (no sort-merge needed); see the
     /// field doc. Mirrors Java `DataSplit#rawConvertible`.
     pub fn raw_convertible(&self) -> bool {

Reply via email to