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 9e7accff feat(table): support selected ranges for rolled dedicated 
files (#592)
9e7accff is described below

commit 9e7accffb463b3367d5a1830bc2059c0b8bd4745
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 23 14:07:15 2026 +0800

    feat(table): support selected ranges for rolled dedicated files (#592)
---
 crates/paimon/src/table/data_evolution_reader.rs | 625 ++++++++++++++++++++++-
 1 file changed, 604 insertions(+), 21 deletions(-)

diff --git a/crates/paimon/src/table/data_evolution_reader.rs 
b/crates/paimon/src/table/data_evolution_reader.rs
index e9cee401..6a792d17 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -32,6 +32,7 @@ use crate::spec::{
 };
 use crate::table::dedicated_format_file_writer::is_blob_file_name;
 use crate::table::schema_manager::SchemaManager;
+use crate::table::source::any_range_overlaps_file;
 use crate::table::{ArrowRecordBatchStream, RESTEnv, RowRange};
 use crate::{DataSplit, Error};
 use arrow_array::{Array, BinaryArray, Int64Array, RecordBatch};
@@ -584,7 +585,13 @@ impl DataEvolutionReader {
                 &prepared_group.files,
             )
             .await?;
-            let source_plan = build_source_plan(&prepared_group, &file_infos, 
&read_type, &blob_descriptor_fields)?;
+            let source_plan = build_source_plan_with_row_id_pushdown(
+                &prepared_group,
+                &file_infos,
+                &read_type,
+                &blob_descriptor_fields,
+                row_ranges.is_some(),
+            )?;
 
             let active_source_indices: Vec<usize> = source_plan
                 .sources
@@ -1128,7 +1135,28 @@ fn open_source_stream(
     blob_as_descriptor: bool,
     anchor_deletion_vector: Option<&DeletionVectorContext>,
 ) -> crate::Result<ArrowRecordBatchStream> {
+    let mut row_ranges = row_ranges;
     if let FieldSource::BlobBunch { bunch, read_fields } = source {
+        let selected_ranges = selected_absolute_row_ranges_for_file(
+            bunch.expected_first_row_id,
+            bunch.expected_row_count,
+            row_ranges.as_deref(),
+            anchor_deletion_vector.map(|context| 
context.deletion_vector.as_ref()),
+        )?;
+        if let Some(selected_ranges) = selected_ranges.as_deref() {
+            let uncovered_ranges =
+                crate::table::source::exclude_row_ranges(selected_ranges, 
bunch.logical_ranges());
+            if !uncovered_ranges.is_empty() {
+                return Err(Error::DataInvalid {
+                    message: format!(
+                        "Blob bunch logical row ranges {:?} do not cover 
effective selected row ranges {selected_ranges:?}; uncovered ranges are 
{uncovered_ranges:?}",
+                        bunch.logical_ranges()
+                    ),
+                    source: None,
+                });
+            }
+        }
+
         // A single sequence group has no fallback work, so keep per-file lazy 
streaming.
         if !bunch.can_read_sequentially() {
             return blob_fallback::read(
@@ -1141,6 +1169,7 @@ fn open_source_stream(
                 anchor_deletion_vector.cloned(),
             );
         }
+        row_ranges = selected_ranges;
     }
 
     let file_reader = DataFileReader::new(
@@ -1168,13 +1197,7 @@ fn open_source_stream(
             )
         }
         FieldSource::BlobBunch { bunch, .. } => {
-            let selected_ranges = selected_absolute_row_ranges_for_file(
-                bunch.expected_first_row_id,
-                bunch.expected_row_count,
-                row_ranges.as_deref(),
-                anchor_deletion_vector.map(|context| 
context.deletion_vector.as_ref()),
-            )?;
-            let files = match selected_ranges.as_deref() {
+            let files = match row_ranges.as_deref() {
                 Some(ranges) => bunch.files_overlapping(ranges)?,
                 None => bunch.files.clone(),
             };
@@ -1189,14 +1212,80 @@ fn open_source_stream(
         }
         FieldSource::VectorBunch {
             bunch, data_fields, ..
-        } => read_bunch_files_stream(
-            file_reader,
-            split,
-            bunch.files.clone(),
-            data_fields.clone(),
-            row_ranges,
-            anchor_deletion_vector.cloned(),
-        ),
+        } => {
+            let anchor = 
crate::table::source::data_evolution_anchor_file(split.data_files())?;
+            let first_row_id = anchor.first_row_id.ok_or_else(|| 
Error::DataInvalid {
+                message: format!(
+                    "Data-evolution anchor file '{}' is missing first_row_id",
+                    anchor.file_name
+                ),
+                source: None,
+            })?;
+            let selected_ranges = selected_absolute_row_ranges_for_file(
+                first_row_id,
+                anchor.row_count,
+                row_ranges.as_deref(),
+                anchor_deletion_vector.map(|context| 
context.deletion_vector.as_ref()),
+            )?;
+            let files = match selected_ranges.as_deref() {
+                Some([]) => Vec::new(),
+                Some(ranges) => {
+                    let covered_ranges = bunch
+                        .files
+                        .iter()
+                        .map(|file| {
+                            let first_row_id =
+                                file.first_row_id.ok_or_else(|| 
Error::DataInvalid {
+                                    message: format!(
+                                        "Vector file '{}' is missing 
first_row_id",
+                                        file.file_name
+                                    ),
+                                    source: None,
+                                })?;
+                            if file.row_count <= 0 {
+                                return Err(Error::DataInvalid {
+                                    message: format!(
+                                        "Vector file '{}' row count must be 
positive, got {}",
+                                        file.file_name, file.row_count
+                                    ),
+                                    source: None,
+                                });
+                            }
+                            let last_row_id = first_row_id
+                                .checked_add(file.row_count - 1)
+                                .ok_or_else(|| Error::DataInvalid {
+                                    message: format!(
+                                        "Vector file '{}' row range overflows 
i64",
+                                        file.file_name
+                                    ),
+                                    source: None,
+                                })?;
+                            Ok(RowRange::new(first_row_id, last_row_id))
+                        })
+                        .collect::<crate::Result<Vec<_>>>()?;
+                    let uncovered_ranges =
+                        crate::table::source::exclude_row_ranges(ranges, 
&covered_ranges);
+                    if !uncovered_ranges.is_empty() {
+                        return Err(Error::DataInvalid {
+                            message: format!(
+                                "Vector bunch does not cover effective 
selected row ranges {uncovered_ranges:?}"
+                            ),
+                            source: None,
+                        });
+                    }
+                    bunch.files_overlapping(ranges)
+                }
+                None => bunch.files.clone(),
+            };
+            read_bunch_files_stream(
+                file_reader,
+                split,
+                files,
+                data_fields.clone(),
+                selected_ranges,
+                anchor_deletion_vector.cloned(),
+            )
+        }
     }
 }
 
@@ -1565,11 +1654,28 @@ struct SourcePlan {
     column_plan: Vec<Option<(usize, usize)>>,
 }
 
+#[cfg(test)]
 fn build_source_plan(
     prepared_group: &PreparedMergeGroup,
     file_infos: &[ResolvedFileInfo],
     read_type: &[DataField],
     blob_descriptor_fields: &HashSet<String>,
+) -> crate::Result<SourcePlan> {
+    build_source_plan_with_row_id_pushdown(
+        prepared_group,
+        file_infos,
+        read_type,
+        blob_descriptor_fields,
+        false,
+    )
+}
+
+fn build_source_plan_with_row_id_pushdown(
+    prepared_group: &PreparedMergeGroup,
+    file_infos: &[ResolvedFileInfo],
+    read_type: &[DataField],
+    blob_descriptor_fields: &HashSet<String>,
+    row_id_pushdown: bool,
 ) -> crate::Result<SourcePlan> {
     let mut sources = Vec::new();
     let mut normal_providers: HashMap<i32, usize> = HashMap::new(); // 
field_id -> source_idx
@@ -1634,7 +1740,8 @@ fn build_source_plan(
                         file.schema_id,
                         format_suffix,
                         normalized.clone(),
-                    ),
+                    )
+                    .with_row_id_pushdown(row_id_pushdown),
                     data_fields: info.data_fields.clone(),
                     read_fields: Vec::new(),
                 });
@@ -1712,7 +1819,7 @@ fn build_source_plan(
         } = source
         {
             bunch.finalize()?;
-            if !read_fields.is_empty() {
+            if !read_fields.is_empty() && !row_id_pushdown {
                 bunch.validate_logical_range()?;
             }
         }
@@ -1723,7 +1830,10 @@ fn build_source_plan(
             bunch, read_fields, ..
         } = source
         {
-            if !read_fields.is_empty() && bunch.row_count() != 
prepared_group.logical_row_count {
+            if !read_fields.is_empty()
+                && !row_id_pushdown
+                && bunch.row_count() != prepared_group.logical_row_count
+            {
                 return Err(Error::DataInvalid {
                     message: format!(
                         "Vector bunch row count {} does not match logical row 
count {}",
@@ -2061,6 +2171,7 @@ struct VectorBunch {
     expected_next_first_row_id: i64,
     latest_max_sequence_number: i64,
     row_count: i64,
+    row_id_pushdown: bool,
 }
 
 impl VectorBunch {
@@ -2080,9 +2191,23 @@ impl VectorBunch {
             expected_next_first_row_id: -1,
             latest_max_sequence_number: -1,
             row_count: 0,
+            row_id_pushdown: false,
         }
     }
 
+    fn with_row_id_pushdown(mut self, row_id_pushdown: bool) -> Self {
+        self.row_id_pushdown = row_id_pushdown;
+        self
+    }
+
+    fn files_overlapping(&self, ranges: &[RowRange]) -> Vec<DataFileMeta> {
+        self.files
+            .iter()
+            .filter(|file| any_range_overlaps_file(ranges, file))
+            .cloned()
+            .collect()
+    }
+
     fn add(&mut self, file: DataFileMeta, normalized_write_cols: &[String]) -> 
crate::Result<()> {
         if !is_vector_store_file_name(&file.file_name) {
             return Err(Error::DataInvalid {
@@ -2117,9 +2242,10 @@ impl VectorBunch {
                                 .to_string(),
                         source: None,
                     });
+                } else {
+                    return Ok(());
                 }
-                return Ok(());
-            } else if first_row_id > self.expected_next_first_row_id {
+            } else if first_row_id > self.expected_next_first_row_id && 
!self.row_id_pushdown {
                 return Err(Error::DataInvalid {
                     message: format!(
                         "Vector file first row id should be continuous, expect 
{} but got {}",
@@ -2886,6 +3012,31 @@ mod tests {
         assert_eq!(bunch.files.len(), 1);
     }
 
+    #[test]
+    fn test_selected_vector_bunch_rejects_partial_higher_sequence_overlap() {
+        let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), 
vec!["emb".to_string()])
+            .with_row_id_pushdown(true);
+        bunch
+            .add(
+                data_file("v-old.vector.parquet", 0, 10, 1, Some(vec!["emb"])),
+                &["emb".to_string()],
+            )
+            .unwrap();
+        let err = bunch
+            .add(
+                data_file("v-new.vector.parquet", 5, 10, 2, Some(vec!["emb"])),
+                &["emb".to_string()],
+            )
+            .unwrap_err();
+
+        assert_eq!(bunch.row_count(), 10);
+        assert_eq!(bunch.files.len(), 1);
+        assert_eq!(bunch.files[0].file_name, "v-old.vector.parquet");
+        assert!(
+            matches!(err, Error::DataInvalid { message, .. } if 
message.contains("overlapping"))
+        );
+    }
+
     #[test]
     fn test_vector_bunch_rejects_row_count_overflow() {
         let mut bunch = VectorBunch::new(15, 0, "parquet".to_string(), 
vec!["emb".to_string()]);
@@ -3267,6 +3418,102 @@ mod tests {
         );
     }
 
+    #[tokio::test]
+    async fn 
test_table_read_accepts_selected_rolled_blob_segment_with_row_ranges() {
+        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 blob_path = bucket_dir.join("blob-part-2.blob");
+        copy_blob_fixture("blob-part-2.blob", &blob_path);
+
+        let file_io = FileIOBuilder::new("file").build().unwrap();
+        let table_schema = TableSchema::new(
+            0,
+            &Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("payload", DataType::Blob(BlobType::new()))
+                .option("data-evolution.enabled", "true")
+                .build()
+                .unwrap(),
+        );
+        let table = Table::new(
+            file_io,
+            Identifier::new("default", "selected_blob_t"),
+            table_path,
+            table_schema,
+            None,
+        );
+
+        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(vec![
+                data_file_meta_with_path(
+                    "data.parquet",
+                    0,
+                    4,
+                    1,
+                    parquet_path.metadata().unwrap().len() as i64,
+                    Some(vec!["id"]),
+                ),
+                data_file_meta_with_path(
+                    "blob-part-2.blob",
+                    2,
+                    2,
+                    1,
+                    blob_path.metadata().unwrap().len() as i64,
+                    Some(vec!["payload"]),
+                ),
+            ])
+            .with_row_ranges(vec![RowRange::new(2, 2)])
+            .build()
+            .unwrap();
+
+        let read = TableRead::new(&table, table.schema().fields().to_vec(), 
Vec::new());
+        let batches = read
+            .to_arrow(std::slice::from_ref(&split))
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+
+        assert_eq!(collect_int_values(&batches, "id"), vec![3]);
+        assert_eq!(
+            collect_binary_values(&batches, "payload"),
+            vec![Some(b"world".to_vec())]
+        );
+
+        let descriptor_table = table.copy_with_options(HashMap::from([(
+            "blob-as-descriptor".to_string(),
+            "true".to_string(),
+        )]));
+        let descriptor_read = TableRead::new(
+            &descriptor_table,
+            descriptor_table.schema().fields().to_vec(),
+            Vec::new(),
+        );
+        let descriptor_batches = descriptor_read
+            .to_arrow(std::slice::from_ref(&split))
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+        let descriptor = collect_binary_values(&descriptor_batches, 
"payload")[0]
+            .clone()
+            .unwrap();
+        let descriptor = BlobDescriptor::deserialize(&descriptor).unwrap();
+        assert!(descriptor.uri().ends_with("blob-part-2.blob"));
+        assert_eq!(descriptor.length(), 5);
+    }
+
     #[tokio::test]
     async fn test_table_read_merges_java_array_blob_file() {
         let tempdir = tempdir().unwrap();
@@ -3815,6 +4062,184 @@ mod tests {
         );
     }
 
+    #[tokio::test]
+    async fn test_selected_blob_fallback_rejects_uncovered_non_deleted_range() 
{
+        use BlobFixtureValue::{Placeholder, Value};
+
+        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 latest_path = bucket_dir.join("blob-latest.blob");
+        let old_path = bucket_dir.join("blob-old.blob");
+        write_blob_file_with_values(&latest_path, &[Placeholder]);
+        write_blob_file_with_values(&old_path, &[Value(b"covered")]);
+
+        let file_io = FileIOBuilder::new("file").build().unwrap();
+        let table_schema = TableSchema::new(
+            0,
+            &Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("payload", DataType::Blob(BlobType::new()))
+                .option("data-evolution.enabled", "true")
+                .build()
+                .unwrap(),
+        );
+        let table = Table::new(
+            file_io.clone(),
+            Identifier::new("default", "selected_blob_gap_t"),
+            table_path,
+            table_schema,
+            None,
+        );
+
+        let deletion_path = format!("{}/index/dv-gap", 
local_file_path(tempdir.path()));
+        let deletion_file = write_test_deletion_file(&file_io, &deletion_path, 
&[1]).await;
+        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(vec![
+                data_file_meta_with_path(
+                    "data.parquet",
+                    0,
+                    4,
+                    1,
+                    parquet_path.metadata().unwrap().len() as i64,
+                    Some(vec!["id"]),
+                ),
+                data_file_meta_with_path(
+                    "blob-latest.blob",
+                    2,
+                    1,
+                    2,
+                    latest_path.metadata().unwrap().len() as i64,
+                    Some(vec!["payload"]),
+                ),
+                data_file_meta_with_path(
+                    "blob-old.blob",
+                    2,
+                    1,
+                    1,
+                    old_path.metadata().unwrap().len() as i64,
+                    Some(vec!["payload"]),
+                ),
+            ])
+            .with_data_deletion_files(vec![Some(deletion_file), None, None])
+            .with_row_ranges(vec![RowRange::new(1, 3)])
+            .build()
+            .unwrap();
+
+        // Row 0 is outside the selection and row 1 is deleted, so only the
+        // uncovered, selected, non-deleted row 3 must make the read fail.
+        for blob_as_descriptor in [false, true] {
+            let mode_table = table.copy_with_options(HashMap::from([(
+                "blob-as-descriptor".to_string(),
+                blob_as_descriptor.to_string(),
+            )]));
+            let read = TableRead::new(
+                &mode_table,
+                mode_table.schema().fields().to_vec(),
+                Vec::new(),
+            );
+            let mut stream = 
read.to_arrow(std::slice::from_ref(&split)).unwrap();
+            let first = stream.try_next().await;
+            assert!(
+                matches!(
+                    &first,
+                    Err(Error::DataInvalid { message, .. })
+                        if message.contains(
+                            "uncovered ranges are [RowRange { from: 3, to: 3 
}]"
+                        )
+                ),
+                "blob_as_descriptor={blob_as_descriptor}: expected uncovered 
selected BLOB range error, got {first:?}"
+            );
+        }
+    }
+
+    #[tokio::test]
+    async fn 
test_selected_blob_range_is_clipped_to_anchor_before_sequential_read() {
+        use BlobFixtureValue::Value;
+
+        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![10, 20])], 
None);
+
+        let blob_path = bucket_dir.join("blob-straddling.blob");
+        write_blob_file_with_values(
+            &blob_path,
+            &[Value(b"outside-anchor"), Value(b"anchor-row-0")],
+        );
+
+        let file_io = FileIOBuilder::new("file").build().unwrap();
+        let table_schema = TableSchema::new(
+            0,
+            &Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("payload", DataType::Blob(BlobType::new()))
+                .option("data-evolution.enabled", "true")
+                .build()
+                .unwrap(),
+        );
+        let table = Table::new(
+            file_io,
+            Identifier::new("default", "selected_blob_anchor_clip_t"),
+            table_path,
+            table_schema,
+            None,
+        );
+
+        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(vec![
+                data_file_meta_with_path(
+                    "data.parquet",
+                    0,
+                    2,
+                    1,
+                    parquet_path.metadata().unwrap().len() as i64,
+                    Some(vec!["id"]),
+                ),
+                data_file_meta_with_path(
+                    "blob-straddling.blob",
+                    -1,
+                    2,
+                    1,
+                    blob_path.metadata().unwrap().len() as i64,
+                    Some(vec!["payload"]),
+                ),
+            ])
+            .with_row_ranges(vec![RowRange::new(-1, 0)])
+            .build()
+            .unwrap();
+
+        let read = TableRead::new(&table, table.schema().fields().to_vec(), 
Vec::new());
+        let mut stream = read.to_arrow(&[split]).unwrap();
+        let first_batch = stream.try_next().await.unwrap().unwrap();
+        assert_eq!(
+            collect_int_values(std::slice::from_ref(&first_batch), "id"),
+            vec![10]
+        );
+        assert_eq!(
+            collect_binary_values(&[first_batch], "payload"),
+            vec![Some(b"anchor-row-0".to_vec())]
+        );
+    }
+
     #[tokio::test]
     async fn test_single_blob_sequence_group_yields_before_opening_next_file() 
{
         let tempdir = tempdir().unwrap();
@@ -4889,6 +5314,164 @@ mod tests {
         );
     }
 
+    #[tokio::test]
+    async fn 
test_read_accepts_selected_rolled_vector_segment_with_row_ranges() {
+        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 normal_path = bucket_dir.join("data.parquet");
+        write_int_parquet_file(&normal_path, vec![("id", vec![1, 2, 3, 4, 5, 
6])], None);
+
+        let vector_path = bucket_dir.join("emb-1.vector.parquet");
+        write_fixed_size_list_parquet(
+            &vector_path,
+            "embedding",
+            2,
+            &[Some(vec![1.0, 1.0]), Some(vec![2.0, 2.0])],
+        );
+        let last_vector_path = bucket_dir.join("emb-3.vector.parquet");
+        write_fixed_size_list_parquet(
+            &last_vector_path,
+            "embedding",
+            2,
+            &[Some(vec![5.0, 5.0]), Some(vec![6.0, 6.0])],
+        );
+
+        let file_io = FileIOBuilder::new("file").build().unwrap();
+        let table_schema = TableSchema::new(
+            0,
+            &Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("embedding", vector_float_type(2))
+                .option("data-evolution.enabled", "true")
+                .build()
+                .unwrap(),
+        );
+        let table = Table::new(
+            file_io,
+            Identifier::new("default", "selected_vector_t"),
+            table_path,
+            table_schema,
+            None,
+        );
+
+        let normal_meta = data_file_meta_with_path(
+            "data.parquet",
+            0,
+            6,
+            1,
+            normal_path.metadata().unwrap().len() as i64,
+            Some(vec!["id"]),
+        );
+        let first_vector_meta = data_file_meta_with_path(
+            "emb-1.vector.parquet",
+            0,
+            2,
+            1,
+            vector_path.metadata().unwrap().len() as i64,
+            Some(vec!["embedding"]),
+        );
+        let last_vector_meta = data_file_meta_with_path(
+            "emb-3.vector.parquet",
+            4,
+            2,
+            1,
+            last_vector_path.metadata().unwrap().len() as i64,
+            Some(vec!["embedding"]),
+        );
+
+        let partial_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(vec![
+                normal_meta.clone(),
+                first_vector_meta.clone(),
+                last_vector_meta.clone(),
+            ])
+            .with_row_ranges(vec![RowRange::new(0, 0), RowRange::new(4, 4)])
+            .build()
+            .unwrap();
+
+        let read = TableRead::new(&table, table.schema().fields().to_vec(), 
Vec::new());
+        let batches = read
+            .to_arrow(&[partial_split])
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+
+        assert_eq!(collect_int_values(&batches, "id"), vec![1, 5]);
+        assert_fixed_size_list(
+            &batches,
+            "embedding",
+            2,
+            &[Some(vec![1.0, 1.0]), Some(vec![5.0, 5.0])],
+        );
+
+        let uncovered_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(vec![
+                normal_meta.clone(),
+                first_vector_meta.clone(),
+                last_vector_meta.clone(),
+            ])
+            .with_row_ranges(vec![RowRange::new(2, 2), RowRange::new(4, 4)])
+            .build()
+            .unwrap();
+        let mut stream = read.to_arrow(&[uncovered_split]).unwrap();
+        let error = stream.try_next().await.unwrap_err();
+        assert!(matches!(error, Error::DataInvalid { message, .. }
+        if message.contains(
+            "does not cover effective selected row ranges [RowRange { from: 2, 
to: 2 }]"
+        )));
+
+        let full_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(vec![
+                normal_meta,
+                first_vector_meta,
+                data_file_meta_with_path(
+                    "missing.vector.parquet",
+                    2,
+                    2,
+                    1,
+                    1,
+                    Some(vec!["embedding"]),
+                ),
+                last_vector_meta,
+            ])
+            .with_row_ranges(vec![RowRange::new(0, 0), RowRange::new(4, 4)])
+            .build()
+            .unwrap();
+        let batches = read
+            .to_arrow(&[full_split])
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+
+        assert_eq!(collect_int_values(&batches, "id"), vec![1, 5]);
+        assert_fixed_size_list(
+            &batches,
+            "embedding",
+            2,
+            &[Some(vec![1.0, 1.0]), Some(vec![5.0, 5.0])],
+        );
+    }
+
     /// (9) row_ranges selecting rows ACROSS a segment boundary -> correct 
subset,
     /// locking in the to_local_row_ranges clip-per-segment behavior.
     ///

Reply via email to