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 8987e452 perf(blob): finalize blob bunch metadata once (#539)
8987e452 is described below

commit 8987e4528287003879af77e5636a19014fe299c9
Author: QuakeWang <[email protected]>
AuthorDate: Mon Jul 20 11:07:44 2026 +0800

    perf(blob): finalize blob bunch metadata once (#539)
---
 crates/paimon/src/table/data_evolution_reader.rs   | 307 +++++++++++++++------
 .../table/data_evolution_reader/blob_fallback.rs   |   3 +-
 2 files changed, 218 insertions(+), 92 deletions(-)

diff --git a/crates/paimon/src/table/data_evolution_reader.rs 
b/crates/paimon/src/table/data_evolution_reader.rs
index e166fea8..14871adf 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -38,9 +38,9 @@ use arrow_array::{Array, BinaryArray, Int64Array, 
RecordBatch};
 use async_stream::try_stream;
 use futures::{StreamExt, TryStreamExt};
 use roaring::RoaringBitmap;
-use std::cmp::Reverse;
-use std::collections::{BTreeMap, HashMap, HashSet};
+use std::collections::{HashMap, HashSet};
 use std::future::Future;
+use std::ops::Range;
 use std::sync::Arc;
 
 /// Whether a file name denotes a dedicated vector-store file 
(`*.vector.<format>`).
@@ -1702,11 +1702,12 @@ fn build_source_plan(
         }
     }
 
-    for source in &sources {
+    for source in &mut sources {
         if let FieldSource::BlobBunch {
             bunch, read_fields, ..
         } = source
         {
+            bunch.finalize()?;
             if !read_fields.is_empty() {
                 bunch.validate_logical_range()?;
             }
@@ -1817,6 +1818,8 @@ struct BlobBunch {
     files: Vec<DataFileMeta>,
     expected_first_row_id: i64,
     expected_row_count: i64,
+    logical_ranges: Option<Vec<RowRange>>,
+    sequence_group_ranges: Vec<Range<usize>>,
 }
 
 impl BlobBunch {
@@ -1825,10 +1828,16 @@ impl BlobBunch {
             files: Vec::new(),
             expected_first_row_id,
             expected_row_count,
+            logical_ranges: None,
+            sequence_group_ranges: Vec::new(),
         }
     }
 
     fn add(&mut self, file: DataFileMeta) -> crate::Result<()> {
+        assert!(
+            self.logical_ranges.is_none(),
+            "Cannot add files to a finalized blob bunch"
+        );
         if !is_blob_file_name(&file.file_name) {
             return Err(Error::DataInvalid {
                 message: "Only blob file can be added to a blob 
bunch.".to_string(),
@@ -1836,7 +1845,7 @@ impl BlobBunch {
             });
         }
 
-        let range = blob_file_row_range(&file)?;
+        blob_file_row_range(&file)?;
         if let Some(first_file) = self.files.first() {
             if file.write_cols != first_file.write_cols {
                 return Err(Error::DataInvalid {
@@ -1847,29 +1856,61 @@ impl BlobBunch {
             }
         }
 
-        for existing in self
-            .files
-            .iter()
-            .filter(|existing| existing.max_sequence_number == 
file.max_sequence_number)
-        {
-            let existing_range = blob_file_row_range(existing)?;
-            if range.overlaps_inclusive(existing_range.from(), 
existing_range.to()) {
-                return Err(Error::DataInvalid {
-                    message: format!(
-                        "Blob files '{}' and '{}' in the same max sequence 
group overlap",
-                        existing.file_name, file.file_name
-                    ),
-                    source: None,
-                });
+        self.files.push(file);
+        Ok(())
+    }
+
+    fn finalize(&mut self) -> crate::Result<()> {
+        assert!(
+            self.logical_ranges.is_none(),
+            "Blob bunch can only be finalized once"
+        );
+
+        self.files.sort_by(|left, right| {
+            right
+                .max_sequence_number
+                .cmp(&left.max_sequence_number)
+                .then_with(|| {
+                    left.first_row_id
+                        .expect("validated blob first_row_id")
+                        .cmp(&right.first_row_id.expect("validated blob 
first_row_id"))
+                })
+        });
+
+        let mut sequence_group_ranges = Vec::new();
+        let mut group_start = 0;
+        while group_start < self.files.len() {
+            let sequence_number = self.files[group_start].max_sequence_number;
+            let mut group_end = group_start + 1;
+            while group_end < self.files.len()
+                && self.files[group_end].max_sequence_number == sequence_number
+            {
+                let previous = &self.files[group_end - 1];
+                let current = &self.files[group_end];
+                let previous_range = blob_file_row_range(previous)?;
+                let current_range = blob_file_row_range(current)?;
+                if current_range.overlaps_inclusive(previous_range.from(), 
previous_range.to()) {
+                    return Err(Error::DataInvalid {
+                        message: format!(
+                            "Blob files '{}' and '{}' in the same max sequence 
group overlap",
+                            previous.file_name, current.file_name
+                        ),
+                        source: None,
+                    });
+                }
+                group_end += 1;
             }
+            sequence_group_ranges.push(group_start..group_end);
+            group_start = group_end;
         }
 
-        let mut ranges = self.logical_ranges();
-        ranges.push(range);
-        let row_count = crate::table::merge_row_ranges(ranges)
-            .iter()
-            .map(RowRange::count)
-            .sum::<i64>();
+        let logical_ranges = crate::table::merge_row_ranges(
+            self.files
+                .iter()
+                .map(blob_file_row_range)
+                .collect::<crate::Result<_>>()?,
+        );
+        let row_count = 
logical_ranges.iter().map(RowRange::count).sum::<i64>();
         if row_count > self.expected_row_count {
             return Err(Error::DataInvalid {
                 message: format!(
@@ -1880,7 +1921,8 @@ impl BlobBunch {
             });
         }
 
-        self.files.push(file);
+        self.logical_ranges = Some(logical_ranges);
+        self.sequence_group_ranges = sequence_group_ranges;
         Ok(())
     }
 
@@ -1888,13 +1930,10 @@ impl BlobBunch {
         self.logical_ranges().iter().map(RowRange::count).sum()
     }
 
-    fn logical_ranges(&self) -> Vec<RowRange> {
-        crate::table::merge_row_ranges(
-            self.files
-                .iter()
-                .map(|file| blob_file_row_range(file).expect("validated blob 
file range"))
-                .collect(),
-        )
+    fn logical_ranges(&self) -> &[RowRange] {
+        self.logical_ranges
+            .as_deref()
+            .expect("blob bunch should be finalized before reading logical 
ranges")
     }
 
     fn expected_range(&self) -> crate::Result<RowRange> {
@@ -1920,7 +1959,7 @@ impl BlobBunch {
     fn validate_logical_range(&self) -> crate::Result<()> {
         let ranges = self.logical_ranges();
         let expected = self.expected_range()?;
-        if ranges.as_slice() != [expected.clone()] {
+        if ranges != [expected.clone()] {
             return Err(Error::DataInvalid {
                 message: format!(
                     "Blob bunch logical row ranges {ranges:?} ({} rows) do not 
match expected range {expected:?}",
@@ -1932,27 +1971,31 @@ impl BlobBunch {
         Ok(())
     }
 
-    fn sequence_groups(&self) -> Vec<Vec<DataFileMeta>> {
-        let mut groups: BTreeMap<Reverse<i64>, Vec<DataFileMeta>> = 
BTreeMap::new();
-        for file in &self.files {
-            groups
-                .entry(Reverse(file.max_sequence_number))
-                .or_default()
-                .push(file.clone());
-        }
-        for files in groups.values_mut() {
-            files.sort_by_key(|file| file.first_row_id.expect("validated blob 
first_row_id"));
+    #[cfg(test)]
+    fn sequence_groups(&self) -> impl Iterator<Item = &[DataFileMeta]> {
+        self.sequence_group_ranges
+            .iter()
+            .map(|range| &self.files[range.clone()])
+    }
+
+    fn into_sequence_groups(self) -> Vec<Vec<DataFileMeta>> {
+        self.logical_ranges
+            .expect("blob bunch should be finalized before reading sequence 
groups");
+
+        let mut files = self.files.into_iter();
+        let mut groups = Vec::with_capacity(self.sequence_group_ranges.len());
+        for range in self.sequence_group_ranges {
+            let group_len = range.len();
+            let group = files.by_ref().take(group_len).collect::<Vec<_>>();
+            debug_assert_eq!(group.len(), group_len);
+            groups.push(group);
         }
-        groups.into_values().collect()
+        debug_assert!(files.next().is_none());
+        groups
     }
 
     fn can_read_sequentially(&self) -> bool {
-        let Some(first_file) = self.files.first() else {
-            return false;
-        };
-        self.files
-            .iter()
-            .all(|file| file.max_sequence_number == 
first_file.max_sequence_number)
+        self.sequence_group_ranges.len() == 1
     }
 
     fn files_overlapping(&self, ranges: &[RowRange]) -> 
crate::Result<Vec<DataFileMeta>> {
@@ -2544,6 +2587,7 @@ mod tests {
         bunch
             .add(data_file("blob-low.blob", 0, 100, 2, Some(vec!["payload"])))
             .unwrap();
+        bunch.finalize().unwrap();
 
         assert_eq!(bunch.row_count(), 100);
         assert_eq!(bunch.files.len(), 2);
@@ -2580,22 +2624,50 @@ mod tests {
     fn test_blob_bunch_groups_sequences_in_descending_order() {
         let mut bunch = BlobBunch::new(0, 1000);
         bunch
-            .add(data_file("blob-low.blob", 0, 100, 2, Some(vec!["payload"])))
+            .add(data_file(
+                "blob-low-late.blob",
+                100,
+                100,
+                2,
+                Some(vec!["payload"]),
+            ))
             .unwrap();
         bunch
             .add(data_file(
-                "blob-high.blob",
+                "blob-high-late.blob",
+                100,
+                100,
+                3,
+                Some(vec!["payload"]),
+            ))
+            .unwrap();
+        bunch
+            .add(data_file(
+                "blob-low-first.blob",
+                0,
+                100,
+                2,
+                Some(vec!["payload"]),
+            ))
+            .unwrap();
+        bunch
+            .add(data_file(
+                "blob-high-first.blob",
                 0,
                 100,
                 3,
                 Some(vec!["payload"]),
             ))
             .unwrap();
+        bunch.finalize().unwrap();
 
-        let groups = bunch.sequence_groups();
-        assert_eq!(groups.len(), 2);
-        assert_eq!(groups[0][0].file_name, "blob-high.blob");
-        assert_eq!(groups[1][0].file_name, "blob-low.blob");
+        assert_eq!(
+            blob_sequence_group_names(&bunch),
+            vec![
+                vec!["blob-high-first.blob", "blob-high-late.blob"],
+                vec!["blob-low-first.blob", "blob-low-late.blob"],
+            ]
+        );
     }
 
     #[test]
@@ -2607,6 +2679,7 @@ mod tests {
         bunch
             .add(data_file("blob2.blob", 50, 150, 2, Some(vec!["payload"])))
             .unwrap();
+        bunch.finalize().unwrap();
 
         assert_eq!(bunch.files.len(), 2);
         assert_eq!(bunch.row_count(), 200);
@@ -2619,9 +2692,10 @@ mod tests {
         bunch
             .add(data_file("blob1.blob", 0, 100, 2, Some(vec!["payload"])))
             .unwrap();
-        let err = bunch
+        bunch
             .add(data_file("blob2.blob", 50, 150, 2, Some(vec!["payload"])))
-            .unwrap_err();
+            .unwrap();
+        let err = bunch.finalize().unwrap_err();
 
         assert!(
             matches!(err, Error::DataInvalid { message, .. } if 
message.contains("same max sequence group"))
@@ -2637,6 +2711,7 @@ mod tests {
         bunch
             .add(data_file("blob2.blob", 150, 100, 2, Some(vec!["payload"])))
             .unwrap();
+        bunch.finalize().unwrap();
         let err = bunch.validate_logical_range().unwrap_err();
 
         assert!(
@@ -2670,6 +2745,7 @@ mod tests {
         let mut mixed_schema = data_file("blob2.blob", 100, 100, 3, 
Some(vec!["payload"]));
         mixed_schema.schema_id = 1;
         bunch.add(mixed_schema).unwrap();
+        bunch.finalize().unwrap();
 
         assert_eq!(bunch.files.len(), 2);
         assert_eq!(bunch.files[0].schema_id, 0);
@@ -2685,10 +2761,10 @@ mod tests {
         bunch
             .add(data_file("blob1.blob", 0, 60, 3, Some(vec!["payload"])))
             .unwrap();
-
-        let err = bunch
+        bunch
             .add(data_file("blob2.blob", 60, 50, 2, Some(vec!["payload"])))
-            .unwrap_err();
+            .unwrap();
+        let err = bunch.finalize().unwrap_err();
 
         assert!(
             matches!(err, Error::DataInvalid { message, .. } if 
message.contains("exceeds the expected"))
@@ -2932,23 +3008,13 @@ mod tests {
 
         match &source_plan.sources[1] {
             FieldSource::BlobBunch { bunch, .. } => {
-                let file_names: Vec<&str> = bunch
-                    .files
-                    .iter()
-                    .map(|file| file.file_name.as_str())
-                    .collect();
                 assert_eq!(
-                    file_names,
+                    blob_sequence_group_names(bunch),
                     vec![
-                        "blob5.blob",
-                        "blob2.blob",
-                        "blob1.blob",
-                        "blob9.blob",
-                        "blob6.blob",
-                        "blob3.blob",
-                        "blob7.blob",
-                        "blob4.blob",
-                        "blob8.blob",
+                        vec!["blob9.blob"],
+                        vec!["blob5.blob", "blob6.blob", "blob7.blob", 
"blob8.blob"],
+                        vec!["blob2.blob", "blob3.blob", "blob4.blob"],
+                        vec!["blob1.blob"],
                     ]
                 );
             }
@@ -2982,6 +3048,64 @@ mod tests {
         assert_eq!(source_plan.column_plan, vec![Some((0, 0)), Some((2, 0))]);
     }
 
+    #[test]
+    fn 
test_build_source_plan_validates_complete_range_only_for_read_blob_source() {
+        let files = vec![
+            data_file("data.parquet", 0, 100, 1, None),
+            data_file("payload.blob", 10, 90, 1, Some(vec!["payload"])),
+        ];
+        let prepared_group = PreparedMergeGroup::new(&files).unwrap();
+        let file_infos = vec![resolved_info(vec![1]), resolved_info(vec![2])];
+        let read_type = vec![DataField::new(
+            1,
+            "id".to_string(),
+            DataType::Int(IntType::new()),
+        )];
+
+        let source_plan =
+            build_source_plan(&prepared_group, &file_infos, &read_type, 
&HashSet::new()).unwrap();
+
+        assert_eq!(source_plan.sources.len(), 2);
+        assert_eq!(source_plan.column_plan, vec![Some((0, 0))]);
+
+        let read_type_with_blob = vec![
+            DataField::new(1, "id".to_string(), DataType::Int(IntType::new())),
+            DataField::new(2, "payload".to_string(), 
DataType::Blob(BlobType::new())),
+        ];
+        let err = build_source_plan(
+            &prepared_group,
+            &file_infos,
+            &read_type_with_blob,
+            &HashSet::new(),
+        )
+        .unwrap_err();
+        assert!(
+            matches!(err, Error::DataInvalid { message, .. } if 
message.contains("logical row ranges"))
+        );
+    }
+
+    #[test]
+    fn test_build_source_plan_rejects_oversized_unread_blob_source() {
+        let files = vec![
+            data_file("data.parquet", 0, 100, 1, None),
+            data_file("payload.blob", 0, 101, 1, Some(vec!["payload"])),
+        ];
+        let prepared_group = PreparedMergeGroup::new(&files).unwrap();
+        let file_infos = vec![resolved_info(vec![1]), resolved_info(vec![2])];
+        let read_type = vec![DataField::new(
+            1,
+            "id".to_string(),
+            DataType::Int(IntType::new()),
+        )];
+
+        let err = build_source_plan(&prepared_group, &file_infos, &read_type, 
&HashSet::new())
+            .unwrap_err();
+
+        assert!(
+            matches!(err, Error::DataInvalid { message, .. } if 
message.contains("exceeds the expected"))
+        );
+    }
+
     #[test]
     fn test_build_source_plan_groups_multiple_blob_columns() {
         let files = vec![
@@ -3024,14 +3148,12 @@ mod tests {
 
         match &source_plan.sources[1] {
             FieldSource::BlobBunch { bunch, .. } => {
-                let file_names: Vec<&str> = bunch
-                    .files
-                    .iter()
-                    .map(|file| file.file_name.as_str())
-                    .collect();
                 assert_eq!(
-                    file_names,
-                    vec!["blob5.blob", "blob9.blob", "blob7.blob", 
"blob8.blob"]
+                    blob_sequence_group_names(bunch),
+                    vec![
+                        vec!["blob9.blob"],
+                        vec!["blob5.blob", "blob7.blob", "blob8.blob"],
+                    ]
                 );
             }
             FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => {
@@ -3041,14 +3163,12 @@ mod tests {
 
         match &source_plan.sources[2] {
             FieldSource::BlobBunch { bunch, .. } => {
-                let file_names: Vec<&str> = bunch
-                    .files
-                    .iter()
-                    .map(|file| file.file_name.as_str())
-                    .collect();
                 assert_eq!(
-                    file_names,
-                    vec!["blob15.blob", "blob19.blob", "blob17.blob", 
"blob18.blob"]
+                    blob_sequence_group_names(bunch),
+                    vec![
+                        vec!["blob19.blob"],
+                        vec!["blob15.blob", "blob17.blob", "blob18.blob"],
+                    ]
                 );
             }
             FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => {
@@ -4732,6 +4852,13 @@ mod tests {
         }
     }
 
+    fn blob_sequence_group_names(bunch: &BlobBunch) -> Vec<Vec<&str>> {
+        bunch
+            .sequence_groups()
+            .map(|files| files.iter().map(|file| 
file.file_name.as_str()).collect())
+            .collect()
+    }
+
     fn data_file(
         file_name: &str,
         first_row_id: i64,
diff --git a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs 
b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
index e34dbca4..105ce4e4 100644
--- a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
+++ b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
@@ -120,7 +120,6 @@ pub(super) fn read(
     let split = split.clone();
 
     Ok(try_stream! {
-        bunch.validate_logical_range()?;
         let expected_range = bunch.expected_range()?;
         let selected_ranges = selected_absolute_row_ranges_for_file(
             bunch.expected_first_row_id,
@@ -133,7 +132,7 @@ pub(super) fn read(
         .unwrap_or_else(|| vec![expected_range]);
 
         let mut sequence_groups = Vec::new();
-        for files in bunch.sequence_groups() {
+        for files in bunch.into_sequence_groups() {
             let mut group = VecDeque::with_capacity(files.len());
             for file in files {
                 let range = blob_file_row_range(&file)?;

Reply via email to