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 0179439d feat(read): load BLOB views lazily under LIMIT (#902)
0179439d is described below

commit 0179439d46d11ae5a256256b9f812069118d4e64
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 18:12:03 2026 +0800

    feat(read): load BLOB views lazily under LIMIT (#902)
---
 crates/paimon/src/table/data_evolution_reader.rs | 47 ++++++++++++++++++++----
 crates/paimon/tests/rest_catalog_test.rs         | 39 +++++++++++++++-----
 2 files changed, 69 insertions(+), 17 deletions(-)

diff --git a/crates/paimon/src/table/data_evolution_reader.rs 
b/crates/paimon/src/table/data_evolution_reader.rs
index 5b6563f1..411ba255 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -294,9 +294,15 @@ impl DataEvolutionReader {
             let descriptor_fields = 
self.descriptor_fields_to_resolve(resolve_blob_views);
             let filter_before_blob_resolution =
                 self.can_filter_before_blob_resolution(resolve_blob_views, 
&descriptor_fields);
-            let blob_view_lookup = self
-                .preload_blob_view_lookup(&splits, 
filter_before_blob_resolution)
-                .await?;
+            // LIMIT reads resolve views only after the surviving batch has
+            // been selected. Eagerly scanning every split here can fail on a
+            // reference beyond the quota and load unrelated upstream blobs.
+            let mut blob_view_lookup = if self.limit.is_some() && 
resolve_blob_views {
+                Some(BlobViewLookup::default())
+            } else {
+                self.preload_blob_view_lookup(&splits, 
filter_before_blob_resolution)
+                    .await?
+            };
             let descriptor_fields = 
self.descriptor_fields_to_resolve(blob_view_lookup.is_some());
             let filter_before_blob_resolution =
                 
self.can_filter_before_blob_resolution(blob_view_lookup.is_some(), 
&descriptor_fields);
@@ -437,7 +443,7 @@ impl DataEvolutionReader {
                                 };
                                 yield self.finish_wide_batch(
                                     batch,
-                                    blob_view_lookup.as_ref(),
+                                    &mut blob_view_lookup,
                                     &descriptor_fields,
                                     filter_before_blob_resolution,
                                     &mut remaining,
@@ -519,7 +525,7 @@ impl DataEvolutionReader {
                             };
                             yield self.finish_wide_batch(
                                 batch,
-                                blob_view_lookup.as_ref(),
+                                &mut blob_view_lookup,
                                 &descriptor_fields,
                                 filter_before_blob_resolution,
                                 &mut remaining,
@@ -589,7 +595,7 @@ impl DataEvolutionReader {
     async fn finish_wide_batch(
         &self,
         batch: RecordBatch,
-        blob_view_lookup: Option<&BlobViewLookup>,
+        blob_view_lookup: &mut Option<BlobViewLookup>,
         descriptor_fields: &HashSet<String>,
         filter_before_blob_resolution: bool,
         remaining: &mut Option<usize>,
@@ -606,7 +612,16 @@ impl DataEvolutionReader {
             return self.project_output(batch);
         }
 
-        batch = self.resolve_blob_view_columns(batch, blob_view_lookup)?;
+        if self.limit.is_some() {
+            if let (Some(lookup), Some(rest_env)) =
+                (blob_view_lookup.as_mut(), self.blob_view_rest_env.clone())
+            {
+                lookup
+                    .load_missing(rest_env, &batch, &self.blob_view_fields)
+                    .await?;
+            }
+        }
+        batch = self.resolve_blob_view_columns(batch, 
blob_view_lookup.as_ref())?;
         let mut batch = if !self.blob_as_descriptor && 
!descriptor_fields.is_empty() {
             resolve_descriptor_columns(
                 batch,
@@ -1214,6 +1229,22 @@ struct BlobViewLookup {
 }
 
 impl BlobViewLookup {
+    async fn load_missing(
+        &mut self,
+        rest_env: RESTEnv,
+        batch: &RecordBatch,
+        blob_view_fields: &HashSet<String>,
+    ) -> crate::Result<()> {
+        let mut view_structs = HashSet::new();
+        collect_blob_view_structs(batch, blob_view_fields, &mut view_structs)?;
+        view_structs.retain(|view| !self.descriptors.contains_key(view));
+        if !view_structs.is_empty() {
+            let loaded = Self::load(rest_env, view_structs).await?;
+            self.descriptors.extend(loaded.descriptors);
+        }
+        Ok(())
+    }
+
     async fn load(rest_env: RESTEnv, view_structs: HashSet<BlobViewStruct>) -> 
crate::Result<Self> {
         if view_structs.is_empty() {
             return Ok(Self::default());
@@ -1341,7 +1372,7 @@ impl BlobViewLookup {
             .map(Option::as_ref)
             .ok_or_else(|| Error::DataInvalid {
                 message: format!(
-                    "BlobViewStruct not found in preloaded cache: 
identifier={}, field_id={}, row_id={}",
+                    "BlobViewStruct not found in lookup cache: identifier={}, 
field_id={}, row_id={}",
                     view_struct.identifier().full_name(),
                     view_struct.field_id(),
                     view_struct.row_id()
diff --git a/crates/paimon/tests/rest_catalog_test.rs 
b/crates/paimon/tests/rest_catalog_test.rs
index bb617d05..baa48cfe 100644
--- a/crates/paimon/tests/rest_catalog_test.rs
+++ b/crates/paimon/tests/rest_catalog_test.rs
@@ -893,7 +893,7 @@ async fn 
test_rest_env_get_table_reuses_catalog_environment() {
 // on Windows elsewhere for the same opendal `fs` StripPrefixError.
 #[cfg(not(windows))]
 #[tokio::test]
-async fn test_blob_view_prescan_filters_invalid_filtered_out_reference() {
+async fn test_blob_view_limit_only_resolves_selected_references() {
     let tmp = tempfile::tempdir().unwrap();
     let warehouse = format!("file://{}", tmp.path().display());
 
@@ -924,7 +924,7 @@ async fn 
test_blob_view_prescan_filters_invalid_filtered_out_reference() {
     .await;
 
     let view_id = Identifier::new("default", "blob_view_target");
-    let view_schema = blob_schema(&[("blob-view-field", "picture")]);
+    let view_schema = blob_schema(&[("blob-view-field", "picture"), 
("read.batch-size", "1")]);
     fs_catalog
         .create_table(&view_id, view_schema.clone(), false)
         .await
@@ -938,18 +938,18 @@ async fn 
test_blob_view_prescan_filters_invalid_filtered_out_reference() {
         .find(|field| field.name() == "picture")
         .unwrap()
         .id();
-    let filtered_out_bad_ref = BlobViewStruct::new(source_id.clone(), 
picture_field_id, 99)
+    let kept_ref = BlobViewStruct::new(source_id.clone(), picture_field_id, 1)
         .serialize()
         .unwrap();
-    let kept_ref = BlobViewStruct::new(source_id.clone(), picture_field_id, 1)
+    let filtered_out_bad_ref = BlobViewStruct::new(source_id.clone(), 
picture_field_id, 99)
         .serialize()
         .unwrap();
     write_batch(
         &view,
         blob_batch(
-            vec![1, 2],
-            vec!["Filtered", "Kept"],
-            vec![filtered_out_bad_ref, kept_ref],
+            vec![1, 2, 3],
+            vec!["Kept", "Repeated", "Filtered"],
+            vec![kept_ref.clone(), kept_ref, filtered_out_bad_ref],
         ),
         "view-writer",
     )
@@ -980,7 +980,7 @@ async fn 
test_blob_view_prescan_filters_invalid_filtered_out_reference() {
 
     let rest_view = rest_catalog.get_table(&view_id).await.unwrap();
     let predicate = PredicateBuilder::new(rest_view.schema().fields())
-        .equal("id", Datum::Int(2))
+        .equal("id", Datum::Int(1))
         .unwrap();
     let mut read_builder = rest_view.new_read_builder();
     read_builder.with_filter(predicate);
@@ -995,7 +995,28 @@ async fn 
test_blob_view_prescan_filters_invalid_filtered_out_reference() {
 
     assert_eq!(
         collect_blob_rows(&batches),
-        vec![(2, "Kept".to_string(), Some(b"bob".to_vec()))]
+        vec![(1, "Kept".to_string(), Some(b"bob".to_vec()))]
+    );
+
+    // LIMIT alone must not resolve the invalid reference in the third row.
+    // The repeated reference is read from the lookup cache in a later batch.
+    let mut limited_builder = rest_view.new_read_builder();
+    limited_builder.with_limit(2);
+    let limited_plan = limited_builder.new_scan().plan().await.unwrap();
+    let limited = limited_builder
+        .new_read()
+        .unwrap()
+        .to_arrow(limited_plan.splits())
+        .unwrap()
+        .try_collect::<Vec<_>>()
+        .await
+        .unwrap();
+    assert_eq!(
+        collect_blob_rows(&limited),
+        vec![
+            (1, "Kept".to_string(), Some(b"bob".to_vec())),
+            (2, "Repeated".to_string(), Some(b"bob".to_vec())),
+        ]
     );
 }
 

Reply via email to