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())),
+ ]
);
}