From 1a0710961bc9422c3277f54c66fc09c8abff3fb3 Mon Sep 17 00:00:00 2001 From: JunRuiLee Date: Tue, 30 Jun 2026 16:07:27 +0800 Subject: [PATCH 1/2] feat(table): read dedicated .vector. files in data-evolution splits Route dedicated `*.vector.` files to a vector column source in the data-evolution read path, so vector columns stored in their own sidecar files are materialized correctly alongside the normal data file. - add `is_vector_store_file_name` detector (`.vector.` segment, case-insensitive) - classify vector files in `normalize_merge_group` (normal -> vector -> blob), exclude them from the merge anchor, and force the column-merge path so a lone vector file is never raw-converted - route vector-typed read fields to the dedicated vector file, falling back to an inline normal-file provider for PR2 compatibility - cover dedicated `.vector.parquet`/`.vector.vortex` reads end-to-end Mirrors upstream Paimon's dedicated vector-store read path. --- .../paimon/src/table/data_evolution_reader.rs | 971 +++++++++++++++++- 1 file changed, 914 insertions(+), 57 deletions(-) diff --git a/crates/paimon/src/table/data_evolution_reader.rs b/crates/paimon/src/table/data_evolution_reader.rs index 487ed9519..4f923f394 100644 --- a/crates/paimon/src/table/data_evolution_reader.rs +++ b/crates/paimon/src/table/data_evolution_reader.rs @@ -33,8 +33,23 @@ use futures::StreamExt; use std::collections::{HashMap, HashSet}; use std::sync::Arc; +/// Whether a file name denotes a dedicated vector-store file (`*.vector.`). +/// Mirrors upstream `VectorType.isVectorStoreFile`: the name contains `.vector.`. +fn is_vector_store_file_name(file_name: &str) -> bool { + file_name.to_ascii_lowercase().contains(".vector.") +} + /// Whether the files in a split can be read independently (no column-wise merge needed). fn is_raw_convertible(files: &[DataFileMeta]) -> bool { + // A split containing a dedicated vector file must go through the column-merge + // path so vector columns are routed to their VectorFile source. Check this + // BEFORE the single-file early-return. + if files + .iter() + .any(|f| is_vector_store_file_name(&f.file_name)) + { + return false; + } if files.len() <= 1 { return true; } @@ -500,6 +515,9 @@ fn open_source_stream( match source { FieldSource::DataFile { file, data_fields, .. + } + | FieldSource::VectorFile { + file, data_fields, .. } => file_reader.read_single_file_stream( split, file.as_ref().clone(), @@ -552,11 +570,13 @@ impl PreparedMergeGroup { let data_files: Vec<&DataFileMeta> = files .iter() - .filter(|file| !is_blob_file_name(&file.file_name)) + .filter(|file| { + !is_blob_file_name(&file.file_name) && !is_vector_store_file_name(&file.file_name) + }) .collect(); if data_files.is_empty() { return Err(Error::DataInvalid { - message: "Field merge split containing .blob files requires at least one non-blob data file".to_string(), + message: "Field merge split with .blob/.vector. files requires at least one normal data file".to_string(), source: None, }); } @@ -655,7 +675,8 @@ fn build_source_plan( blob_descriptor_fields: &HashSet, ) -> crate::Result { let mut sources = Vec::new(); - let mut normal_source_indices: HashMap = HashMap::new(); + let mut normal_providers: HashMap = HashMap::new(); // field_id -> source_idx + let mut vector_providers: HashMap = HashMap::new(); // field_id -> source_idx let mut blob_source_indices: HashMap = HashMap::new(); let mut expected_blob_row_count: Option = None; @@ -688,6 +709,26 @@ fn build_source_plan( .blob_bunch_mut() .unwrap() .add(file.clone())?; + } else if is_vector_store_file_name(&file.file_name) { + // A vector file is a column provider only; unlike a normal data file it does + // NOT update `expected_blob_row_count` (it must not anchor a following blob's + // row count). + let source_idx = sources.len(); + sources.push(FieldSource::VectorFile { + file: Box::new(file.clone()), + data_fields: info.data_fields.clone(), + read_fields: Vec::new(), + }); + for &field_id in &info.field_ids { + if vector_providers.insert(field_id, source_idx).is_some() { + return Err(Error::DataInvalid { + message: format!( + "Vector field id {field_id} is provided by more than one .vector. file" + ), + source: None, + }); + } + } } else { expected_blob_row_count = Some(file.row_count); let source_idx = sources.len(); @@ -696,7 +737,10 @@ fn build_source_plan( data_fields: info.data_fields.clone(), read_fields: Vec::new(), }); - normal_source_indices.insert(file_idx, source_idx); + for &field_id in &info.field_ids { + // first normal file that carries the id wins (preserve existing semantics) + normal_providers.entry(field_id).or_insert(source_idx); + } } } @@ -706,13 +750,16 @@ fn build_source_plan( && !blob_descriptor_fields.contains(field.name()) { blob_source_indices.get(&field.id()).copied() + } else if matches!(field.data_type(), DataType::Vector(_)) { + // Prefer the dedicated .vector. file; fall back to a normal data file + // (PR 2 inline-vector compatibility path). + vector_providers + .get(&field.id()) + .copied() + .or_else(|| normal_providers.get(&field.id()).copied()) } else { - select_normal_provider( - &prepared_group.files, - file_infos, - &normal_source_indices, - field.id(), - ) + // Non-vector fields never read from a .vector. file. + normal_providers.get(&field.id()).copied() }; if let Some(source_idx) = source_idx { @@ -747,25 +794,6 @@ fn build_source_plan( }) } -fn select_normal_provider( - files: &[DataFileMeta], - file_infos: &[ResolvedFileInfo], - normal_source_indices: &HashMap, - field_id: i32, -) -> Option { - files.iter().enumerate().find_map(|(file_idx, file)| { - if is_blob_file_name(&file.file_name) { - return None; - } - - file_infos[file_idx] - .field_ids - .contains(&field_id) - .then(|| normal_source_indices.get(&file_idx).copied()) - .flatten() - }) -} - fn resolve_blob_field_id(file: &DataFileMeta, info: &ResolvedFileInfo) -> crate::Result { if info.field_ids.len() != 1 { return Err(Error::DataInvalid { @@ -788,6 +816,11 @@ enum FieldSource { data_fields: Option>, read_fields: Vec, }, + VectorFile { + file: Box, + data_fields: Option>, + read_fields: Vec, + }, BlobBunch { bunch: BlobBunch, data_fields: Option>, @@ -799,6 +832,7 @@ impl FieldSource { fn read_fields(&self) -> &[DataField] { match self { FieldSource::DataFile { read_fields, .. } + | FieldSource::VectorFile { read_fields, .. } | FieldSource::BlobBunch { read_fields, .. } => read_fields, } } @@ -806,6 +840,7 @@ impl FieldSource { fn add_read_field(&mut self, field: DataField) -> usize { let read_fields = match self { FieldSource::DataFile { read_fields, .. } + | FieldSource::VectorFile { read_fields, .. } | FieldSource::BlobBunch { read_fields, .. } => read_fields, }; if let Some(offset) = read_fields @@ -822,7 +857,7 @@ impl FieldSource { fn blob_bunch_mut(&mut self) -> Option<&mut BlobBunch> { match self { FieldSource::BlobBunch { bunch, .. } => Some(bunch), - FieldSource::DataFile { .. } => None, + FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => None, } } } @@ -937,44 +972,53 @@ impl BlobBunch { } fn normalize_merge_group(files: Vec) -> crate::Result> { - let mut data_files = Vec::new(); + let mut normal_files = Vec::new(); + let mut vector_files = Vec::new(); let mut blob_files = Vec::new(); for file in files { if is_blob_file_name(&file.file_name) { blob_files.push(file); + } else if is_vector_store_file_name(&file.file_name) { + vector_files.push(file); } else { - data_files.push(file); + normal_files.push(file); } } - data_files.sort_by_key(|f| std::cmp::Reverse(f.max_sequence_number)); - if let Some(first) = data_files.first() { - let first_row_id = first.first_row_id.ok_or_else(|| Error::DataInvalid { - message: "All data files in a field merge split should have first_row_id".to_string(), + normal_files.sort_by_key(|f| std::cmp::Reverse(f.max_sequence_number)); + vector_files.sort_by_key(|f| std::cmp::Reverse(f.max_sequence_number)); + + // Normal + vector files all share the anchor's row range. Validate together, + // using the first normal file as the reference when present, else the first + // vector file (a vector-only group has no anchor and is rejected later in + // PreparedMergeGroup::new). + let mut range_ref: Option<(i64, i64)> = None; + for file in normal_files.iter().chain(vector_files.iter()) { + let first_row_id = file.first_row_id.ok_or_else(|| Error::DataInvalid { + message: "All data/vector files in a field merge split should have first_row_id" + .to_string(), source: None, })?; - let first_row_count = first.row_count; - for file in data_files.iter().skip(1) { - if file.first_row_id != Some(first_row_id) || file.row_count != first_row_count { - return Err(Error::DataInvalid { - message: - "All data files in a field merge split should have the same row id range." - .to_string(), - source: None, - }); + match range_ref { + None => range_ref = Some((first_row_id, file.row_count)), + Some((ref_first, ref_count)) => { + if first_row_id != ref_first || file.row_count != ref_count { + return Err(Error::DataInvalid { + message: "All data and vector files in a field merge split should have the same row id range.".to_string(), + source: None, + }); + } } } } blob_files.sort_by(|left, right| { - let left_first_row_id = left.first_row_id.unwrap_or(i64::MIN); - let right_first_row_id = right.first_row_id.unwrap_or(i64::MIN); - left_first_row_id - .cmp(&right_first_row_id) + let l = left.first_row_id.unwrap_or(i64::MIN); + let r = right.first_row_id.unwrap_or(i64::MIN); + l.cmp(&r) .then_with(|| right.max_sequence_number.cmp(&left.max_sequence_number)) }); - if blob_files.iter().any(|file| file.first_row_id.is_none()) { return Err(Error::DataInvalid { message: "All blob files in a field merge split should have first_row_id".to_string(), @@ -982,8 +1026,10 @@ fn normalize_merge_group(files: Vec) -> crate::Result vector -> blob regardless of sequence, yielding [d1, v1, b1]. + let files = vec![ + data_file("v1.vector.parquet", 0, 10, 5, Some(vec!["emb"])), + data_file("b1.blob", 0, 1, 1, Some(vec!["payload"])), + data_file("d1.parquet", 0, 10, 1, Some(vec!["id"])), + ]; + let normalized = normalize_merge_group(files).unwrap(); + let names: Vec<&str> = normalized.iter().map(|f| f.file_name.as_str()).collect(); + // normal first, then vector, then blob + assert_eq!(names, vec!["d1.parquet", "v1.vector.parquet", "b1.blob"]); + } + + #[test] + fn test_normalize_merge_group_rejects_vector_row_range_mismatch() { + let files = vec![ + data_file("d1.parquet", 0, 10, 1, Some(vec!["id"])), + data_file("v1.vector.parquet", 0, 8, 1, Some(vec!["emb"])), // row_count mismatch + ]; + let err = normalize_merge_group(files); + assert!(matches!(err, Err(Error::DataInvalid { .. }))); + } + #[test] fn test_blob_bunch_ignores_same_first_row_id_with_lower_sequence() { let mut bunch = BlobBunch::new(1000); @@ -1077,6 +1191,31 @@ mod tests { assert_eq!(bunch.files[0].file_name, "blob-high.blob"); } + #[test] + fn test_is_vector_store_file_name() { + assert!(is_vector_store_file_name("data-1.vector.parquet")); + assert!(is_vector_store_file_name("data-1.vector.vortex")); + assert!(is_vector_store_file_name("PART.VECTOR.PARQUET")); // case-insensitive + assert!(!is_vector_store_file_name("data-1.parquet")); + assert!(!is_vector_store_file_name("data-1.blob")); + assert!(!is_vector_store_file_name("x.vectorstuff")); // not the ".vector." segment + } + + #[test] + fn test_is_raw_convertible_false_for_single_vector_file() { + // A lone vector file must NOT be raw-convertible (would bypass merge routing). + let files = vec![data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"]))]; + assert!(!is_raw_convertible(&files)); + } + + #[test] + fn test_prepared_merge_group_rejects_vector_only_split() { + // No normal anchor file -> DataInvalid. + let files = vec![data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"]))]; + let err = PreparedMergeGroup::new(&files); + assert!(matches!(err, Err(Error::DataInvalid { .. }))); + } + #[test] fn test_blob_bunch_rejects_same_first_row_id_with_higher_sequence() { let mut bunch = BlobBunch::new(1000); @@ -1228,7 +1367,9 @@ mod tests { vec!["blob5.blob", "blob9.blob", "blob7.blob", "blob8.blob"] ); } - FieldSource::DataFile { .. } => panic!("expected blob bunch source"), + FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => { + panic!("expected blob bunch source") + } } } @@ -1308,7 +1449,9 @@ mod tests { vec!["blob5.blob", "blob9.blob", "blob7.blob", "blob8.blob"] ); } - FieldSource::DataFile { .. } => panic!("expected blob bunch source"), + FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => { + panic!("expected blob bunch source") + } } match &source_plan.sources[2] { @@ -1323,7 +1466,9 @@ mod tests { vec!["blob15.blob", "blob19.blob", "blob17.blob", "blob18.blob"] ); } - FieldSource::DataFile { .. } => panic!("expected blob bunch source"), + FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => { + panic!("expected blob bunch source") + } } } @@ -1524,6 +1669,718 @@ mod tests { ); } + fn write_fixed_size_list_parquet( + path: &std::path::Path, + col: &str, + dim: i32, + rows: &[Option>], + ) { + use arrow_array::builder::{FixedSizeListBuilder, Float32Builder}; + use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema}; + use parquet::arrow::ArrowWriter; + use std::fs::File; + + let mut builder = FixedSizeListBuilder::new(Float32Builder::new(), dim).with_field( + Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)), + ); + for row in rows { + match row { + Some(vals) => { + assert_eq!(vals.len() as i32, dim); + for v in vals { + builder.values().append_value(*v); + } + builder.append(true); + } + None => { + for _ in 0..dim { + builder.values().append_value(0.0); + } + builder.append(false); + } + } + } + let array = builder.finish(); + let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new( + col, + ArrowDataType::FixedSizeList( + Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)), + dim, + ), + true, + )])); + let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(array)]).unwrap(); + let file = File::create(path).unwrap(); + let mut writer = ArrowWriter::try_new(file, schema, None).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + } + + /// Build a `VECTOR` column type whose element is nullable, matching the + /// arrow `FixedSizeList(element: Float32 nullable)` produced by the writer helper. + fn vector_float_type(dim: u32) -> DataType { + DataType::Vector(VectorType::try_new(true, dim, DataType::Float(FloatType::new())).unwrap()) + } + + /// Locate the embedding column, downcast to `FixedSizeListArray`, and assert the + /// per-row validity bitmap and child `Float32` values across all batches. + fn assert_fixed_size_list( + batches: &[RecordBatch], + column_name: &str, + expected_dim: i32, + expected: &[Option>], + ) { + let mut row = 0usize; + for batch in batches { + let idx = batch.schema().index_of(column_name).unwrap(); + let list = batch + .column(idx) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(list.value_length(), expected_dim); + for i in 0..list.len() { + let want = &expected[row]; + match want { + Some(vals) => { + assert!(list.is_valid(i), "row {row} expected non-null"); + let child = list.value(i); + let floats = child.as_any().downcast_ref::().unwrap(); + let got: Vec = (0..floats.len()).map(|j| floats.value(j)).collect(); + assert_eq!(&got, vals, "row {row} value mismatch"); + } + None => { + assert!(list.is_null(i), "row {row} expected null"); + } + } + row += 1; + } + } + assert_eq!(row, expected.len(), "row count mismatch"); + } + + /// (1) Provider priority: the normal data file ALSO advertises the embedding write_col, + /// but the dedicated `.vector.parquet` file must win. + #[tokio::test] + async fn test_read_dedicated_vector_parquet_file_with_provider_priority() { + 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(); + + // Normal data file carries id AND a (wrong) inline embedding to prove priority. + let normal_path = bucket_dir.join("data.parquet"); + write_int_parquet_file(&normal_path, vec![("id", vec![1, 2, 3])], None); + + // Dedicated vector file: row1=[1,2], row2=null, row3=[3,4]. + let vector_path = bucket_dir.join("data.vector.parquet"); + write_fixed_size_list_parquet( + &vector_path, + "embedding", + 2, + &[Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.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", "vec_priority_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![ + // Normal file advertises BOTH id and embedding write_cols. + data_file_meta_with_path( + "data.parquet", + 0, + 3, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id", "embedding"]), + ), + data_file_meta_with_path( + "data.vector.parquet", + 0, + 3, + 1, + vector_path.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + ]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3]); + // Value MUST come from the .vector. file (vector-provider priority). + assert_fixed_size_list( + &batches, + "embedding", + 2, + &[Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.0])], + ); + } + + /// (2) Same shape but the dedicated vector file is `.vector.vortex`. + #[cfg(feature = "vortex")] + #[tokio::test] + async fn test_read_dedicated_vector_vortex_file() { + use crate::arrow::format::create_format_writer; + use arrow_array::builder::{FixedSizeListBuilder, Float32Builder}; + use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema}; + + 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])], None); + + // Write data.vector.vortex via the format writer (dispatches on the .vortex suffix). + let vector_path = bucket_dir.join("data.vector.vortex"); + let file_io = FileIOBuilder::new("file").build().unwrap(); + { + let mut builder = FixedSizeListBuilder::new(Float32Builder::new(), 2).with_field( + Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)), + ); + for row in [Some([1.0_f32, 2.0]), None, Some([3.0, 4.0])] { + match row { + Some(vals) => { + for v in vals { + builder.values().append_value(v); + } + builder.append(true); + } + None => { + builder.values().append_value(0.0); + builder.values().append_value(0.0); + builder.append(false); + } + } + } + let array = builder.finish(); + let arrow_schema = Arc::new(ArrowSchema::new(vec![ArrowField::new( + "embedding", + ArrowDataType::FixedSizeList( + Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)), + 2, + ), + true, + )])); + let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(array)]).unwrap(); + let output = file_io.new_output(&local_file_path(&vector_path)).unwrap(); + let mut writer = create_format_writer(&output, arrow_schema, "zstd", 1, None) + .await + .unwrap(); + writer.write(&batch).await.unwrap(); + writer.close().await.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", "vec_vortex_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, + 3, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + ), + data_file_meta_with_path( + "data.vector.vortex", + 0, + 3, + 1, + vector_path.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + ]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3]); + assert_fixed_size_list( + &batches, + "embedding", + 2, + &[Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.0])], + ); + } + + /// (3) Multiple vector columns living in ONE `.vector.parquet` file; both must + /// route to the same VectorFile source and materialize. + #[tokio::test] + async fn test_read_dedicated_vector_file_multiple_columns() { + use arrow_array::builder::{FixedSizeListBuilder, Float32Builder}; + use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema}; + use parquet::arrow::ArrowWriter; + use std::fs::File; + + 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])], None); + + // One vector file with two FixedSizeList columns: emb1 (dim 2), emb2 (dim 3). + let vector_path = bucket_dir.join("data.vector.parquet"); + { + let elem = || Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)); + let mut b1 = FixedSizeListBuilder::new(Float32Builder::new(), 2).with_field(elem()); + let mut b2 = FixedSizeListBuilder::new(Float32Builder::new(), 3).with_field(elem()); + // emb1: [1,2], null, [5,6] + for row in [Some(vec![1.0_f32, 2.0]), None, Some(vec![5.0, 6.0])] { + match row { + Some(v) => { + for x in v { + b1.values().append_value(x); + } + b1.append(true); + } + None => { + b1.values().append_value(0.0); + b1.values().append_value(0.0); + b1.append(false); + } + } + } + // emb2: [7,8,9], [1,1,1], null + for row in [ + Some(vec![7.0_f32, 8.0, 9.0]), + Some(vec![1.0, 1.0, 1.0]), + None, + ] { + match row { + Some(v) => { + for x in v { + b2.values().append_value(x); + } + b2.append(true); + } + None => { + for _ in 0..3 { + b2.values().append_value(0.0); + } + b2.append(false); + } + } + } + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("emb1", ArrowDataType::FixedSizeList(elem(), 2), true), + ArrowField::new("emb2", ArrowDataType::FixedSizeList(elem(), 3), true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![Arc::new(b1.finish()), Arc::new(b2.finish())], + ) + .unwrap(); + let file = File::create(&vector_path).unwrap(); + let mut writer = ArrowWriter::try_new(file, schema, None).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + } + + let file_io = FileIOBuilder::new("file").build().unwrap(); + let table_schema = TableSchema::new( + 0, + &Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("emb1", vector_float_type(2)) + .column("emb2", vector_float_type(3)) + .option("data-evolution.enabled", "true") + .build() + .unwrap(), + ); + let table = Table::new( + file_io, + Identifier::new("default", "vec_multi_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, + 3, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + ), + data_file_meta_with_path( + "data.vector.parquet", + 0, + 3, + 1, + vector_path.metadata().unwrap().len() as i64, + Some(vec!["emb1", "emb2"]), + ), + ]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3]); + assert_fixed_size_list( + &batches, + "emb1", + 2, + &[Some(vec![1.0, 2.0]), None, Some(vec![5.0, 6.0])], + ); + assert_fixed_size_list( + &batches, + "emb2", + 3, + &[Some(vec![7.0, 8.0, 9.0]), Some(vec![1.0, 1.0, 1.0]), None], + ); + } + + /// (4) Inline fallback: embedding lives in the normal parquet, NO `.vector.` file + /// present. Routing must fall back to the normal provider (PR 2 compatibility). + #[tokio::test] + async fn test_inline_vector_fallback_still_reads() { + 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(); + + // Single normal parquet holding id + an inline FixedSizeList embedding. + let normal_path = bucket_dir.join("data.parquet"); + { + use arrow_array::builder::{FixedSizeListBuilder, Float32Builder}; + use arrow_schema::{ + DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema, + }; + use parquet::arrow::ArrowWriter; + use std::fs::File; + + let mut emb = FixedSizeListBuilder::new(Float32Builder::new(), 2).with_field(Arc::new( + ArrowField::new("element", ArrowDataType::Float32, true), + )); + for row in [Some([1.0_f32, 2.0]), None, Some([3.0, 4.0])] { + match row { + Some(vals) => { + for v in vals { + emb.values().append_value(v); + } + emb.append(true); + } + None => { + emb.values().append_value(0.0); + emb.values().append_value(0.0); + emb.append(false); + } + } + } + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("id", ArrowDataType::Int32, false), + ArrowField::new( + "embedding", + ArrowDataType::FixedSizeList( + Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)), + 2, + ), + true, + ), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![1, 2, 3])), + Arc::new(emb.finish()), + ], + ) + .unwrap(); + let file = File::create(&normal_path).unwrap(); + let mut writer = ArrowWriter::try_new(file, schema, None).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + } + + 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", "vec_inline_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, + 3, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id", "embedding"]), + )]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3]); + assert_fixed_size_list( + &batches, + "embedding", + 2, + &[Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.0])], + ); + } + + /// (5) A `.vector.` file is present, but a non-vector field (`id`) must still be + /// read from the normal file, never mis-selected from the vector file. + #[tokio::test] + async fn test_non_vector_field_ignores_vector_file() { + 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![10, 20, 30])], None); + + let vector_path = bucket_dir.join("data.vector.parquet"); + write_fixed_size_list_parquet( + &vector_path, + "embedding", + 2, + &[Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.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", "vec_nonvec_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, + 3, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + ), + data_file_meta_with_path( + "data.vector.parquet", + 0, + 3, + 1, + vector_path.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + ]) + .build() + .unwrap(); + + // Project only the non-vector `id` field. + let id_field = table + .schema() + .fields() + .iter() + .find(|f| f.name() == "id") + .unwrap() + .clone(); + let read = TableRead::new(&table, vec![id_field], Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![10, 20, 30]); + } + + /// (6) Row-range mismatch: normal file row_count=3 but `.vector.parquet` row_count=2 + /// must surface as DataInvalid. + #[tokio::test] + async fn test_read_vector_file_row_range_mismatch_errors() { + 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])], None); + + let vector_path = bucket_dir.join("data.vector.parquet"); + write_fixed_size_list_parquet( + &vector_path, + "embedding", + 2, + &[Some(vec![1.0, 2.0]), Some(vec![3.0, 4.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", "vec_mismatch_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, + 3, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + ), + data_file_meta_with_path( + "data.vector.parquet", + 0, + 2, // row_count mismatch vs the normal file's 3 + 1, + vector_path.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + ]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let result = read.to_arrow(&[split]); + let collected = match result { + Ok(stream) => stream.try_collect::>().await, + Err(e) => Err(e), + }; + assert!( + matches!(collected, Err(Error::DataInvalid { .. })), + "expected DataInvalid, got {collected:?}" + ); + } + fn resolved_info(field_ids: Vec) -> ResolvedFileInfo { ResolvedFileInfo { field_ids, From dede55096a4be0a19641411095f135043039eec1 Mon Sep 17 00:00:00 2001 From: JunRuiLee Date: Tue, 30 Jun 2026 16:07:56 +0800 Subject: [PATCH 2/2] feat(table): reassemble rolling (multi-segment) vector files When a vector column is rolled into multiple `*.vector.` segments (upstream `DedicatedFormatRollingFileWriter` uses an independent `VECTOR_TARGET_FILE_SIZE`), the read path must reassemble them. This mirrors upstream Paimon's `VectorFileBunch` non-pushdown semantics. - add `VectorBunch` (modeled on `BlobBunch`): aggregates segments of one logical vector source, keyed by (schema_id, format suffix, normalized write cols), with continuity/dedup guards in `add` (same-first_row_id lower-seq ignore, overlap lower-seq ignore, gap error, row-count overflow error, key-identity mismatch error) - sort vector segments by first_row_id asc / max_sequence_number desc in `normalize_merge_group`, and stop validating rolled segments (slices) against the normal-file row range - compute per-file `normalized_write_cols` in `load_file_infos` (write cols sorted by field position; missing write_cols -> DataInvalid), so differently-ordered raw write cols that normalize equal aggregate into one bunch - aggregate same-key segments into one `FieldSource::VectorBunch` in `build_source_plan`, with a bunch-granularity duplicate-field-id guard and a bunch row-count check; a single vector file is simply a bunch of length 1 - read the bunch via the same sequential-concat path as the blob bunch; each segment carries its real first_row_id/row_count, so `to_local_row_ranges` clips absolute row ranges per segment - cover rolled reassembly and cross-boundary row_ranges end-to-end Out of scope: write side, blob fallback-placeholder, predicate/vector pushdown, upstream rowIdPushDown branch. --- .../paimon/src/table/data_evolution_reader.rs | 1126 +++++++++++++++-- 1 file changed, 1008 insertions(+), 118 deletions(-) diff --git a/crates/paimon/src/table/data_evolution_reader.rs b/crates/paimon/src/table/data_evolution_reader.rs index 4f923f394..8f4a5cb6a 100644 --- a/crates/paimon/src/table/data_evolution_reader.rs +++ b/crates/paimon/src/table/data_evolution_reader.rs @@ -42,11 +42,11 @@ fn is_vector_store_file_name(file_name: &str) -> bool { /// Whether the files in a split can be read independently (no column-wise merge needed). fn is_raw_convertible(files: &[DataFileMeta]) -> bool { // A split containing a dedicated vector file must go through the column-merge - // path so vector columns are routed to their VectorFile source. Check this + // path so vector columns are routed to their VectorBunch source. Check this // BEFORE the single-file early-return. if files .iter() - .any(|f| is_vector_store_file_name(&f.file_name)) + .any(|file| is_vector_store_file_name(&file.file_name)) { return false; } @@ -515,9 +515,6 @@ fn open_source_stream( match source { FieldSource::DataFile { file, data_fields, .. - } - | FieldSource::VectorFile { - file, data_fields, .. } => file_reader.read_single_file_stream( split, file.as_ref().clone(), @@ -527,27 +524,48 @@ fn open_source_stream( ), FieldSource::BlobBunch { bunch, data_fields, .. - } => { - let split = split.clone(); - let files = bunch.files.clone(); - let data_fields = data_fields.clone(); - Ok(try_stream! { - for file in files { - let mut stream = file_reader.read_single_file_stream( - &split, - file, - data_fields.clone(), - None, - row_ranges.clone(), - )?; - while let Some(batch) = stream.next().await { - yield batch?; - } - } + } => read_bunch_files_stream( + file_reader, + split, + bunch.files.clone(), + data_fields.clone(), + row_ranges, + ), + FieldSource::VectorBunch { + bunch, data_fields, .. + } => read_bunch_files_stream( + file_reader, + split, + bunch.files.clone(), + data_fields.clone(), + row_ranges, + ), + } +} + +fn read_bunch_files_stream( + file_reader: DataFileReader, + split: &DataSplit, + files: Vec, + data_fields: Option>, + row_ranges: Option>, +) -> crate::Result { + let split = split.clone(); + Ok(try_stream! { + for file in files { + let mut stream = file_reader.read_single_file_stream( + &split, + file, + data_fields.clone(), + None, + row_ranges.clone(), + )?; + while let Some(batch) = stream.next().await { + yield batch?; } - .boxed()) } } + .boxed()) } #[derive(Debug, Clone)] @@ -611,6 +629,7 @@ impl PreparedMergeGroup { struct ResolvedFileInfo { field_ids: Vec, data_fields: Option>, + normalized_write_cols: Option>, } async fn load_file_infos( @@ -622,19 +641,34 @@ async fn load_file_infos( let mut infos = Vec::with_capacity(files.len()); for file in files { + let (field_ids, data_fields, effective_fields_owned); if file.schema_id == table_schema_id { - infos.push(ResolvedFileInfo { - field_ids: resolve_field_ids(file, table_fields)?, - data_fields: None, - }); + field_ids = resolve_field_ids(file, table_fields)?; + data_fields = None; + effective_fields_owned = None; } else { let data_schema = schema_manager.schema(file.schema_id).await?; - let data_fields = data_schema.fields().to_vec(); - infos.push(ResolvedFileInfo { - field_ids: resolve_field_ids(file, &data_fields)?, - data_fields: Some(data_fields), - }); + let fields = data_schema.fields().to_vec(); + field_ids = resolve_field_ids(file, &fields)?; + data_fields = Some(fields.clone()); + effective_fields_owned = Some(fields); } + + let normalized_write_cols = if is_vector_store_file_name(&file.file_name) { + let effective_fields: &[DataField] = match effective_fields_owned.as_deref() { + Some(fields) => fields, + None => table_fields, + }; + Some(normalize_vector_write_cols(file, effective_fields)?) + } else { + None + }; + + infos.push(ResolvedFileInfo { + field_ids, + data_fields, + normalized_write_cols, + }); } Ok(infos) @@ -662,6 +696,50 @@ fn resolve_field_ids(file: &DataFileMeta, fields: &[DataField]) -> crate::Result } } +/// Lowercased final filename extension, used as a vector bunch's format identifier. +/// `"data.vector.parquet" -> "parquet"`, `"emb-1.vector.vortex" -> "vortex"`. +fn vector_format_suffix(file_name: &str) -> String { + file_name + .rsplit('.') + .next() + .unwrap_or("") + .to_ascii_lowercase() +} + +/// Normalize a vector file's write columns into a stable key component: the write +/// column names sorted by their field position in the file's effective row type. +/// A `.vector.` file with no `write_cols` is ambiguous and rejected. An unknown +/// column name is rejected. Raw orderings that differ but normalize equal compare equal. +fn normalize_vector_write_cols( + file: &DataFileMeta, + fields: &[DataField], +) -> crate::Result> { + let write_cols = file.write_cols.as_ref().ok_or_else(|| Error::DataInvalid { + message: format!("Vector file '{}' must declare write_cols", file.file_name), + source: None, + })?; + + let mut indexed: Vec<(usize, String)> = write_cols + .iter() + .map(|name| { + fields + .iter() + .position(|field| field.name() == name) + .map(|pos| (pos, name.clone())) + .ok_or_else(|| Error::DataInvalid { + message: format!( + "Failed to resolve vector write column '{}' in file '{}'", + name, file.file_name + ), + source: None, + }) + }) + .collect::>()?; + + indexed.sort_by_key(|(pos, _)| *pos); + Ok(indexed.into_iter().map(|(_, name)| name).collect()) +} + #[derive(Debug, Clone)] struct SourcePlan { sources: Vec, @@ -676,7 +754,8 @@ fn build_source_plan( ) -> crate::Result { let mut sources = Vec::new(); let mut normal_providers: HashMap = HashMap::new(); // field_id -> source_idx - let mut vector_providers: HashMap = HashMap::new(); // field_id -> source_idx + let mut vector_field_providers: HashMap = HashMap::new(); // field_id -> source_idx + let mut vector_bunch_indices: HashMap<(i64, String, Vec), usize> = HashMap::new(); let mut blob_source_indices: HashMap = HashMap::new(); let mut expected_blob_row_count: Option = None; @@ -712,21 +791,60 @@ fn build_source_plan( } else if is_vector_store_file_name(&file.file_name) { // A vector file is a column provider only; unlike a normal data file it does // NOT update `expected_blob_row_count` (it must not anchor a following blob's - // row count). - let source_idx = sources.len(); - sources.push(FieldSource::VectorFile { - file: Box::new(file.clone()), - data_fields: info.data_fields.clone(), - read_fields: Vec::new(), - }); - for &field_id in &info.field_ids { - if vector_providers.insert(field_id, source_idx).is_some() { - return Err(Error::DataInvalid { + // row count). Segments sharing the same (schema_id, format, normalized + // write cols) key aggregate into one bunch. + let normalized = + info.normalized_write_cols + .clone() + .ok_or_else(|| Error::DataInvalid { message: format!( - "Vector field id {field_id} is provided by more than one .vector. file" + "Vector file '{}' is missing normalized write columns", + file.file_name ), source: None, - }); + })?; + let format_suffix = vector_format_suffix(&file.file_name); + let key = (file.schema_id, format_suffix.clone(), normalized.clone()); + + let source_idx = if let Some(&existing_idx) = vector_bunch_indices.get(&key) { + existing_idx + } else { + let source_idx = sources.len(); + sources.push(FieldSource::VectorBunch { + bunch: VectorBunch::new( + prepared_group.logical_row_count, + file.schema_id, + format_suffix, + normalized.clone(), + ), + data_fields: info.data_fields.clone(), + read_fields: Vec::new(), + }); + vector_bunch_indices.insert(key, source_idx); + source_idx + }; + + sources[source_idx] + .vector_bunch_mut() + .unwrap() + .add(file.clone(), &normalized)?; + + for &field_id in &info.field_ids { + match vector_field_providers.get(&field_id) { + // Same bunch aggregating another segment: fine. + Some(&existing_idx) if existing_idx == source_idx => {} + // Different bunch key advertising the same field id: ambiguous. + Some(_) => { + return Err(Error::DataInvalid { + message: format!( + "Vector field id {field_id} is provided by more than one vector bunch" + ), + source: None, + }); + } + None => { + vector_field_providers.insert(field_id, source_idx); + } } } } else { @@ -751,9 +869,9 @@ fn build_source_plan( { blob_source_indices.get(&field.id()).copied() } else if matches!(field.data_type(), DataType::Vector(_)) { - // Prefer the dedicated .vector. file; fall back to a normal data file + // Prefer the dedicated .vector. bunch; fall back to a normal data file // (PR 2 inline-vector compatibility path). - vector_providers + vector_field_providers .get(&field.id()) .copied() .or_else(|| normal_providers.get(&field.id()).copied()) @@ -788,6 +906,24 @@ fn build_source_plan( } } + for source in &sources { + if let FieldSource::VectorBunch { + bunch, read_fields, .. + } = source + { + if !read_fields.is_empty() && bunch.row_count() != prepared_group.logical_row_count { + return Err(Error::DataInvalid { + message: format!( + "Vector bunch row count {} does not match logical row count {}", + bunch.row_count(), + prepared_group.logical_row_count + ), + source: None, + }); + } + } + } + Ok(SourcePlan { sources, column_plan, @@ -816,8 +952,8 @@ enum FieldSource { data_fields: Option>, read_fields: Vec, }, - VectorFile { - file: Box, + VectorBunch { + bunch: VectorBunch, data_fields: Option>, read_fields: Vec, }, @@ -832,7 +968,7 @@ impl FieldSource { fn read_fields(&self) -> &[DataField] { match self { FieldSource::DataFile { read_fields, .. } - | FieldSource::VectorFile { read_fields, .. } + | FieldSource::VectorBunch { read_fields, .. } | FieldSource::BlobBunch { read_fields, .. } => read_fields, } } @@ -840,7 +976,7 @@ impl FieldSource { fn add_read_field(&mut self, field: DataField) -> usize { let read_fields = match self { FieldSource::DataFile { read_fields, .. } - | FieldSource::VectorFile { read_fields, .. } + | FieldSource::VectorBunch { read_fields, .. } | FieldSource::BlobBunch { read_fields, .. } => read_fields, }; if let Some(offset) = read_fields @@ -857,7 +993,14 @@ impl FieldSource { fn blob_bunch_mut(&mut self) -> Option<&mut BlobBunch> { match self { FieldSource::BlobBunch { bunch, .. } => Some(bunch), - FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => None, + FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => None, + } + } + + fn vector_bunch_mut(&mut self) -> Option<&mut VectorBunch> { + match self { + FieldSource::VectorBunch { bunch, .. } => Some(bunch), + FieldSource::DataFile { .. } | FieldSource::BlobBunch { .. } => None, } } } @@ -971,6 +1114,137 @@ impl BlobBunch { } } +/// Aggregates rolled `.vector.` segments belonging to one logical vector +/// source, mirroring upstream `VectorFileBunch` non-pushdown semantics. Unlike +/// `BlobBunch`, the expected row count is taken directly from the prepared group's +/// logical row count (vectors sit before blobs and never anchor a blob's row count). +/// +/// `normalize_merge_group` is responsible for ordering segments; `add` assumes sorted +/// input and enforces continuity/dedup. +#[derive(Debug, Clone)] +struct VectorBunch { + files: Vec, + schema_id: i64, + format_suffix: String, + normalized_write_cols: Vec, + expected_row_count: i64, + latest_first_row_id: i64, + expected_next_first_row_id: i64, + latest_max_sequence_number: i64, + row_count: i64, +} + +impl VectorBunch { + fn new( + expected_row_count: i64, + schema_id: i64, + format_suffix: String, + normalized_write_cols: Vec, + ) -> Self { + Self { + files: Vec::new(), + schema_id, + format_suffix, + normalized_write_cols, + expected_row_count, + latest_first_row_id: -1, + expected_next_first_row_id: -1, + latest_max_sequence_number: -1, + row_count: 0, + } + } + + 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 { + message: "Only vector file can be added to a vector bunch.".to_string(), + source: None, + }); + } + + 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 first_row_id == self.latest_first_row_id { + if file.max_sequence_number >= self.latest_max_sequence_number { + return Err(Error::DataInvalid { + message: + "Vector file with same first row id should have decreasing sequence number." + .to_string(), + source: None, + }); + } + return Ok(()); + } + + if !self.files.is_empty() { + if first_row_id < self.expected_next_first_row_id { + if file.max_sequence_number >= self.latest_max_sequence_number { + return Err(Error::DataInvalid { + message: + "Vector file with overlapping row id should have decreasing sequence number." + .to_string(), + source: None, + }); + } + return Ok(()); + } else if first_row_id > self.expected_next_first_row_id { + return Err(Error::DataInvalid { + message: format!( + "Vector file first row id should be continuous, expect {} but got {}", + self.expected_next_first_row_id, first_row_id + ), + source: None, + }); + } + } + + // Defensive key-identity check against the bunch's key (not raw write_cols). + if file.schema_id != self.schema_id { + return Err(Error::DataInvalid { + message: "All files in a vector bunch should have the same schema id.".to_string(), + source: None, + }); + } + if vector_format_suffix(&file.file_name) != self.format_suffix { + return Err(Error::DataInvalid { + message: "All files in a vector bunch should have the same format.".to_string(), + source: None, + }); + } + if normalized_write_cols != self.normalized_write_cols.as_slice() { + return Err(Error::DataInvalid { + message: + "All files in a vector bunch should have the same normalized write columns." + .to_string(), + source: None, + }); + } + + self.row_count += file.row_count; + if self.row_count > self.expected_row_count { + return Err(Error::DataInvalid { + message: format!( + "Vector files row count {} exceed the expected {}", + self.row_count, self.expected_row_count + ), + source: None, + }); + } + self.latest_max_sequence_number = file.max_sequence_number; + self.latest_first_row_id = first_row_id; + self.expected_next_first_row_id = first_row_id + file.row_count; + self.files.push(file); + Ok(()) + } + + fn row_count(&self) -> i64 { + self.row_count + } +} + fn normalize_merge_group(files: Vec) -> crate::Result> { let mut normal_files = Vec::new(); let mut vector_files = Vec::new(); @@ -987,17 +1261,28 @@ fn normalize_merge_group(files: Vec) -> crate::Result = None; - for file in normal_files.iter().chain(vector_files.iter()) { + for file in normal_files.iter() { let first_row_id = file.first_row_id.ok_or_else(|| Error::DataInvalid { - message: "All data/vector files in a field merge split should have first_row_id" - .to_string(), + message: "All data files in a field merge split should have first_row_id".to_string(), source: None, })?; match range_ref { @@ -1005,7 +1290,7 @@ fn normalize_merge_group(files: Vec) -> crate::Result { if first_row_id != ref_first || file.row_count != ref_count { return Err(Error::DataInvalid { - message: "All data and vector files in a field merge split should have the same row id range.".to_string(), + message: "All data files in a field merge split should have the same row id range.".to_string(), source: None, }); } @@ -1078,41 +1363,139 @@ mod tests { use test_utils::{local_file_path, write_int_parquet_file}; #[test] - fn test_build_source_plan_rejects_duplicate_vector_field_id() { - // Two .vector. files both advertising the same vector field id (1) must be - // rejected rather than one silently overwriting the other. + fn test_build_source_plan_aggregates_same_key_vector_segments() { + // Two contiguous vector segments, same key -> ONE VectorBunch, files in sorted order. let files = vec![ - data_file("v1.vector.parquet", 0, 10, 2, Some(vec!["emb"])), - data_file("v2.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + data_file("d1.parquet", 0, 20, 1, Some(vec!["id"])), + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["emb"])), ]; let prepared_group = PreparedMergeGroup { files: files.clone(), - logical_row_count: 10, + logical_row_count: 20, first_row_id: 0, }; let file_infos = vec![ ResolvedFileInfo { field_ids: vec![1], data_fields: None, + normalized_write_cols: None, + }, + ResolvedFileInfo { + field_ids: vec![2], + data_fields: None, + normalized_write_cols: Some(vec!["emb".to_string()]), + }, + ResolvedFileInfo { + field_ids: vec![2], + data_fields: None, + normalized_write_cols: Some(vec!["emb".to_string()]), }, + ]; + let read_type = vec![ + DataField::new(1, "id".to_string(), DataType::Int(IntType::new())), + DataField::new(2, "emb".to_string(), vector_float_type(2)), + ]; + let plan = + build_source_plan(&prepared_group, &file_infos, &read_type, &HashSet::new()).unwrap(); + + // sources: [DataFile(d1), VectorBunch(v1,v2)] + assert_eq!(plan.sources.len(), 2); + assert_eq!(plan.column_plan, vec![Some((0, 0)), Some((1, 0))]); + match &plan.sources[1] { + FieldSource::VectorBunch { bunch, .. } => { + let names: Vec<&str> = bunch.files.iter().map(|f| f.file_name.as_str()).collect(); + assert_eq!(names, vec!["v1.vector.parquet", "v2.vector.parquet"]); + } + _ => panic!("expected vector bunch source"), + } + } + + #[test] + fn test_build_source_plan_aggregates_differently_ordered_write_cols() { + // Two segments with multiple vector cols whose RAW write_cols differ in order but + // normalize to the same key -> one bunch (#5b). field 2 = "a", field 3 = "b". + let files = vec![ + data_file("d1.parquet", 0, 20, 1, Some(vec!["id"])), + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["a", "b"])), + data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["b", "a"])), + ]; + let prepared_group = PreparedMergeGroup { + files: files.clone(), + logical_row_count: 20, + first_row_id: 0, + }; + // Both segments normalize to ["a","b"] (field-position order). + let file_infos = vec![ ResolvedFileInfo { field_ids: vec![1], data_fields: None, + normalized_write_cols: None, + }, + ResolvedFileInfo { + field_ids: vec![2, 3], + data_fields: None, + normalized_write_cols: Some(vec!["a".to_string(), "b".to_string()]), + }, + ResolvedFileInfo { + field_ids: vec![2, 3], + data_fields: None, + normalized_write_cols: Some(vec!["a".to_string(), "b".to_string()]), }, ]; - let read_type = vec![DataField::new( - 1, - "emb".to_string(), - DataType::Vector( - VectorType::try_new(true, 2, DataType::Float(FloatType::new())).unwrap(), - ), - )]; - let err = build_source_plan( - &prepared_group, - &file_infos, - &read_type, - &std::collections::HashSet::new(), - ); + let read_type = vec![ + DataField::new(1, "id".to_string(), DataType::Int(IntType::new())), + DataField::new(2, "a".to_string(), vector_float_type(2)), + DataField::new(3, "b".to_string(), vector_float_type(2)), + ]; + let plan = + build_source_plan(&prepared_group, &file_infos, &read_type, &HashSet::new()).unwrap(); + // One vector bunch holding both segments; both vector columns map to it. + assert_eq!(plan.sources.len(), 2); + match &plan.sources[1] { + FieldSource::VectorBunch { bunch, .. } => assert_eq!(bunch.files.len(), 2), + _ => panic!("expected vector bunch source"), + } + assert_eq!(plan.column_plan[1].map(|(s, _)| s), Some(1)); + assert_eq!(plan.column_plan[2].map(|(s, _)| s), Some(1)); + } + + #[test] + fn test_build_source_plan_rejects_field_id_across_two_bunch_keys() { + // Same field id 2 advertised by two DIFFERENT bunch keys (different write col sets) -> error (#6). + let files = vec![ + data_file("d1.parquet", 0, 10, 1, Some(vec!["id"])), + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + data_file("v2.vector.parquet", 0, 10, 2, Some(vec!["emb", "other"])), + ]; + let prepared_group = PreparedMergeGroup { + files: files.clone(), + logical_row_count: 10, + first_row_id: 0, + }; + let file_infos = vec![ + ResolvedFileInfo { + field_ids: vec![1], + data_fields: None, + normalized_write_cols: None, + }, + ResolvedFileInfo { + field_ids: vec![2], + data_fields: None, + normalized_write_cols: Some(vec!["emb".to_string()]), + }, + ResolvedFileInfo { + field_ids: vec![2, 3], + data_fields: None, + normalized_write_cols: Some(vec!["emb".to_string(), "other".to_string()]), + }, + ]; + let read_type = vec![ + DataField::new(1, "id".to_string(), DataType::Int(IntType::new())), + DataField::new(2, "emb".to_string(), vector_float_type(2)), + DataField::new(3, "other".to_string(), vector_float_type(2)), + ]; + let err = build_source_plan(&prepared_group, &file_infos, &read_type, &HashSet::new()); assert!(matches!(err, Err(Error::DataInvalid { .. }))); } @@ -1161,10 +1544,56 @@ mod tests { } #[test] - fn test_normalize_merge_group_rejects_vector_row_range_mismatch() { + fn test_normalize_merge_group_accepts_rolled_vectors_with_differing_ranges() { + // Rolled vector segments are slices with differing row ranges; they must NOT be + // rejected against the normal anchor's full range (inverts the old reject test). + let files = vec![ + data_file("d1.parquet", 0, 20, 1, Some(vec!["id"])), + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["emb"])), + ]; + let normalized = normalize_merge_group(files).unwrap(); + let names: Vec<&str> = normalized.iter().map(|f| f.file_name.as_str()).collect(); + assert_eq!( + names, + vec!["d1.parquet", "v1.vector.parquet", "v2.vector.parquet"] + ); + } + + #[test] + fn test_normalize_merge_group_sorts_multi_segment_vectors() { + // Vectors out of order: must sort by first_row_id asc, then max_seq desc, + // and land after normal, before blob. + let files = vec![ + data_file("b1.blob", 0, 1, 1, Some(vec!["payload"])), + data_file("v-mid.vector.parquet", 10, 10, 1, Some(vec!["emb"])), + data_file("d1.parquet", 0, 30, 1, Some(vec!["id"])), + data_file("v-late-low.vector.parquet", 20, 10, 1, Some(vec!["emb"])), + data_file("v-late-high.vector.parquet", 20, 10, 5, Some(vec!["emb"])), + data_file("v-early.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + ]; + let normalized = normalize_merge_group(files).unwrap(); + let names: Vec<&str> = normalized.iter().map(|f| f.file_name.as_str()).collect(); + assert_eq!( + names, + vec![ + "d1.parquet", + "v-early.vector.parquet", + "v-mid.vector.parquet", + "v-late-high.vector.parquet", // same first_row_id 20, higher seq first + "v-late-low.vector.parquet", + "b1.blob", + ] + ); + } + + #[test] + fn test_normalize_merge_group_requires_first_row_id_on_vector_files() { + let mut vector_no_rid = data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])); + vector_no_rid.first_row_id = None; let files = vec![ data_file("d1.parquet", 0, 10, 1, Some(vec!["id"])), - data_file("v1.vector.parquet", 0, 8, 1, Some(vec!["emb"])), // row_count mismatch + vector_no_rid, ]; let err = normalize_merge_group(files); assert!(matches!(err, Err(Error::DataInvalid { .. }))); @@ -1255,67 +1684,265 @@ mod tests { } #[test] - fn test_blob_bunch_rejects_non_continuous_first_row_id() { - let mut bunch = BlobBunch::new(1000); + fn test_blob_bunch_rejects_non_continuous_first_row_id() { + let mut bunch = BlobBunch::new(1000); + bunch + .add(data_file("blob1.blob", 0, 100, 3, Some(vec!["payload"]))) + .unwrap(); + + let err = bunch + .add(data_file("blob2.blob", 150, 100, 2, Some(vec!["payload"]))) + .unwrap_err(); + + assert!( + matches!(err, Error::DataInvalid { message, .. } if message.contains("continuous")) + ); + } + + #[test] + fn test_blob_bunch_rejects_mixed_write_columns() { + let mut bunch = BlobBunch::new(200); + bunch + .add(data_file("blob1.blob", 0, 100, 3, Some(vec!["payload"]))) + .unwrap(); + + let err = bunch + .add(data_file("blob2.blob", 100, 100, 2, Some(vec!["payload2"]))) + .unwrap_err(); + + assert!( + matches!(err, Error::DataInvalid { message, .. } if message.contains("same write columns")) + ); + } + + #[test] + fn test_blob_bunch_rejects_mixed_schema_ids() { + let mut bunch = BlobBunch::new(200); + bunch + .add(data_file("blob1.blob", 0, 100, 3, Some(vec!["payload"]))) + .unwrap(); + + let mut mixed_schema = data_file("blob2.blob", 100, 100, 2, Some(vec!["payload"])); + mixed_schema.schema_id = 1; + let err = bunch.add(mixed_schema).unwrap_err(); + + assert!( + matches!(err, Error::DataInvalid { message, .. } if message.contains("same schema id")) + ); + } + + #[test] + fn test_blob_bunch_rejects_row_count_exceeding_expected() { + let mut bunch = BlobBunch::new(100); + bunch + .add(data_file("blob1.blob", 0, 60, 3, Some(vec!["payload"]))) + .unwrap(); + + let err = bunch + .add(data_file("blob2.blob", 60, 50, 2, Some(vec!["payload"]))) + .unwrap_err(); + + assert!( + matches!(err, Error::DataInvalid { message, .. } if message.contains("exceed the expected")) + ); + } + + #[test] + fn test_vector_bunch_aggregates_contiguous_segments() { + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); + bunch + .add( + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + bunch + .add( + data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + bunch + .add( + data_file("v3.vector.parquet", 20, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + assert_eq!(bunch.row_count(), 30); + let names: Vec<&str> = bunch.files.iter().map(|f| f.file_name.as_str()).collect(); + assert_eq!( + names, + vec![ + "v1.vector.parquet", + "v2.vector.parquet", + "v3.vector.parquet" + ] + ); + } + + #[test] + fn test_vector_bunch_rejects_gap() { + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); bunch - .add(data_file("blob1.blob", 0, 100, 3, Some(vec!["payload"]))) + .add( + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap(); - + // first_row_id 15 > expected_next 10 -> gap let err = bunch - .add(data_file("blob2.blob", 150, 100, 2, Some(vec!["payload"]))) + .add( + data_file("v2.vector.parquet", 15, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap_err(); - assert!( matches!(err, Error::DataInvalid { message, .. } if message.contains("continuous")) ); } #[test] - fn test_blob_bunch_rejects_mixed_write_columns() { - let mut bunch = BlobBunch::new(200); + fn test_vector_bunch_ignores_same_first_row_id_lower_seq() { + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); bunch - .add(data_file("blob1.blob", 0, 100, 3, Some(vec!["payload"]))) + .add( + data_file("v-high.vector.parquet", 0, 10, 3, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + // same first_row_id, strictly lower seq -> ignored (dedup), no row_count contribution + bunch + .add( + data_file("v-low.vector.parquet", 0, 10, 2, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap(); + assert_eq!(bunch.row_count(), 10); + assert_eq!(bunch.files.len(), 1); + assert_eq!(bunch.files[0].file_name, "v-high.vector.parquet"); + } + #[test] + fn test_vector_bunch_rejects_same_first_row_id_higher_seq() { + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); + bunch + .add( + data_file("v-low.vector.parquet", 0, 10, 2, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); let err = bunch - .add(data_file("blob2.blob", 100, 100, 2, Some(vec!["payload2"]))) + .add( + data_file("v-high.vector.parquet", 0, 10, 3, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap_err(); - - assert!( - matches!(err, Error::DataInvalid { message, .. } if message.contains("same write columns")) - ); + assert!(matches!(err, Error::DataInvalid { .. })); } #[test] - fn test_blob_bunch_rejects_mixed_schema_ids() { - let mut bunch = BlobBunch::new(200); + fn test_vector_bunch_ignores_overlapping_lower_seq() { + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); bunch - .add(data_file("blob1.blob", 0, 100, 3, Some(vec!["payload"]))) + .add( + data_file("v1.vector.parquet", 0, 10, 3, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap(); + // first_row_id 5 < expected_next 10 -> overlap; lower seq -> ignored + bunch + .add( + data_file("v2.vector.parquet", 5, 10, 2, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + assert_eq!(bunch.row_count(), 10); + assert_eq!(bunch.files.len(), 1); + } - let mut mixed_schema = data_file("blob2.blob", 100, 100, 2, Some(vec!["payload"])); - mixed_schema.schema_id = 1; - let err = bunch.add(mixed_schema).unwrap_err(); - - assert!( - matches!(err, Error::DataInvalid { message, .. } if message.contains("same schema id")) - ); + #[test] + fn test_vector_bunch_rejects_row_count_overflow() { + let mut bunch = VectorBunch::new(15, 0, "parquet".to_string(), vec!["emb".to_string()]); + bunch + .add( + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + let err = bunch + .add( + data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap_err(); + assert!(matches!(err, Error::DataInvalid { message, .. } if message.contains("exceed"))); } #[test] - fn test_blob_bunch_rejects_row_count_exceeding_expected() { - let mut bunch = BlobBunch::new(100); + fn test_vector_bunch_rejects_key_identity_mismatch() { + // schema_id mismatch + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); bunch - .add(data_file("blob1.blob", 0, 60, 3, Some(vec!["payload"]))) + .add( + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap(); + let mut wrong_schema = data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["emb"])); + wrong_schema.schema_id = 99; + let err = bunch.add(wrong_schema, &["emb".to_string()]).unwrap_err(); + assert!(matches!(err, Error::DataInvalid { .. })); + + // format_suffix mismatch + let mut bunch2 = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); + bunch2 + .add( + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + let err2 = bunch2 + .add( + data_file("v2.vector.vortex", 10, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap_err(); + assert!(matches!(err2, Error::DataInvalid { .. })); + + // normalized_write_cols mismatch + let mut bunch3 = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); + bunch3 + .add( + data_file("v1.vector.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) + .unwrap(); + let err3 = bunch3 + .add( + data_file("v2.vector.parquet", 10, 10, 1, Some(vec!["other"])), + &["other".to_string()], + ) + .unwrap_err(); + assert!(matches!(err3, Error::DataInvalid { .. })); + } + #[test] + fn test_vector_bunch_rejects_non_vector_file() { + let mut bunch = VectorBunch::new(30, 0, "parquet".to_string(), vec!["emb".to_string()]); let err = bunch - .add(data_file("blob2.blob", 60, 50, 2, Some(vec!["payload"]))) + .add( + data_file("v1.parquet", 0, 10, 1, Some(vec!["emb"])), + &["emb".to_string()], + ) .unwrap_err(); + assert!(matches!(err, Error::DataInvalid { .. })); + } - assert!( - matches!(err, Error::DataInvalid { message, .. } if message.contains("exceed the expected")) - ); + #[test] + fn test_vector_format_suffix() { + assert_eq!(vector_format_suffix("data.vector.parquet"), "parquet"); + assert_eq!(vector_format_suffix("emb-1.vector.vortex"), "vortex"); + assert_eq!(vector_format_suffix("X.VECTOR.PARQUET"), "parquet"); } #[test] @@ -1367,7 +1994,7 @@ mod tests { vec!["blob5.blob", "blob9.blob", "blob7.blob", "blob8.blob"] ); } - FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => { + FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => { panic!("expected blob bunch source") } } @@ -1449,7 +2076,7 @@ mod tests { vec!["blob5.blob", "blob9.blob", "blob7.blob", "blob8.blob"] ); } - FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => { + FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => { panic!("expected blob bunch source") } } @@ -1466,7 +2093,7 @@ mod tests { vec!["blob15.blob", "blob19.blob", "blob17.blob", "blob18.blob"] ); } - FieldSource::DataFile { .. } | FieldSource::VectorFile { .. } => { + FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => { panic!("expected blob bunch source") } } @@ -1722,6 +2349,35 @@ mod tests { DataType::Vector(VectorType::try_new(true, dim, DataType::Float(FloatType::new())).unwrap()) } + #[test] + fn test_normalize_vector_write_cols_sorts_by_field_position() { + let fields = vec![ + DataField::new(1, "id".to_string(), DataType::Int(IntType::new())), + DataField::new(2, "a".to_string(), vector_float_type(2)), + DataField::new(3, "b".to_string(), vector_float_type(2)), + ]; + // raw write_cols listed b, a -> normalized must be a, b (field-position order) + let file = data_file("v.vector.parquet", 0, 10, 1, Some(vec!["b", "a"])); + let normalized = normalize_vector_write_cols(&file, &fields).unwrap(); + assert_eq!(normalized, vec!["a".to_string(), "b".to_string()]); + } + + #[test] + fn test_normalize_vector_write_cols_rejects_missing_write_cols() { + let fields = vec![DataField::new(2, "a".to_string(), vector_float_type(2))]; + let file = data_file("v.vector.parquet", 0, 10, 1, None); + let err = normalize_vector_write_cols(&file, &fields).unwrap_err(); + assert!(matches!(err, Error::DataInvalid { .. })); + } + + #[test] + fn test_normalize_vector_write_cols_rejects_unknown_column() { + let fields = vec![DataField::new(2, "a".to_string(), vector_float_type(2))]; + let file = data_file("v.vector.parquet", 0, 10, 1, Some(vec!["ghost"])); + let err = normalize_vector_write_cols(&file, &fields).unwrap_err(); + assert!(matches!(err, Error::DataInvalid { .. })); + } + /// Locate the embedding column, downcast to `FixedSizeListArray`, and assert the /// per-row validity bitmap and child `Float32` values across all batches. fn assert_fixed_size_list( @@ -1894,7 +2550,7 @@ mod tests { )])); let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(array)]).unwrap(); let output = file_io.new_output(&local_file_path(&vector_path)).unwrap(); - let mut writer = create_format_writer(&output, arrow_schema, "zstd", 1, None) + let mut writer = create_format_writer(&output, arrow_schema, "zstd", 1, None, None) .await .unwrap(); writer.write(&batch).await.unwrap(); @@ -1963,7 +2619,7 @@ mod tests { } /// (3) Multiple vector columns living in ONE `.vector.parquet` file; both must - /// route to the same VectorFile source and materialize. + /// route to the same VectorBunch source and materialize. #[tokio::test] async fn test_read_dedicated_vector_file_multiple_columns() { use arrow_array::builder::{FixedSizeListBuilder, Float32Builder}; @@ -2304,6 +2960,239 @@ mod tests { assert_eq!(collect_int_values(&batches, "id"), vec![10, 20, 30]); } + /// (8) normal data.parquet (id) + 3 rolled .vector.parquet segments (embedding, + /// contiguous row ranges) reassemble into one column with values in correct order. + #[tokio::test] + async fn test_read_rolled_vector_segments_reassemble() { + 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(); + + // Normal data file: id 1..=6 (6 rows total). + let normal_path = bucket_dir.join("data.parquet"); + write_int_parquet_file(&normal_path, vec![("id", vec![1, 2, 3, 4, 5, 6])], None); + + // Three rolled vector segments, 2 rows each, contiguous first_row_ids 0,2,4. + let seg1 = bucket_dir.join("emb-1.vector.parquet"); + write_fixed_size_list_parquet( + &seg1, + "embedding", + 2, + &[Some(vec![1.0, 1.0]), Some(vec![2.0, 2.0])], + ); + let seg2 = bucket_dir.join("emb-2.vector.parquet"); + write_fixed_size_list_parquet(&seg2, "embedding", 2, &[Some(vec![3.0, 3.0]), None]); + let seg3 = bucket_dir.join("emb-3.vector.parquet"); + write_fixed_size_list_parquet( + &seg3, + "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", "vec_rolled_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, + 6, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + ), + data_file_meta_with_path( + "emb-1.vector.parquet", + 0, + 2, + 1, + seg1.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + data_file_meta_with_path( + "emb-2.vector.parquet", + 2, + 2, + 1, + seg2.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + data_file_meta_with_path( + "emb-3.vector.parquet", + 4, + 2, + 1, + seg3.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + ]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3, 4, 5, 6]); + assert_fixed_size_list( + &batches, + "embedding", + 2, + &[ + Some(vec![1.0, 1.0]), + Some(vec![2.0, 2.0]), + Some(vec![3.0, 3.0]), + None, + Some(vec![5.0, 5.0]), + Some(vec![6.0, 6.0]), + ], + ); + } + + /// (9) row_ranges selecting rows ACROSS a segment boundary -> correct subset, + /// locking in the to_local_row_ranges clip-per-segment behavior. + /// + /// `RowRange::new` is inclusive on both ends (see source::RowRange::count), so the + /// absolute window [1, 3] selects rows at index 1,2,3 (ids 2,3,4). Row 1 lives in + /// segment emb-1 [0,2) and rows 2,3 live in emb-2 [2,4), so the window straddles the + /// emb-1/emb-2 boundary and must be clipped per segment via `to_local_row_ranges`. + #[tokio::test] + async fn test_read_rolled_vector_segments_with_cross_boundary_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(); + + // Normal data file: id 1..=6 (6 rows total). + let normal_path = bucket_dir.join("data.parquet"); + write_int_parquet_file(&normal_path, vec![("id", vec![1, 2, 3, 4, 5, 6])], None); + + // Three rolled vector segments, 2 rows each, contiguous first_row_ids 0,2,4. + let seg1 = bucket_dir.join("emb-1.vector.parquet"); + write_fixed_size_list_parquet( + &seg1, + "embedding", + 2, + &[Some(vec![1.0, 1.0]), Some(vec![2.0, 2.0])], + ); + let seg2 = bucket_dir.join("emb-2.vector.parquet"); + write_fixed_size_list_parquet(&seg2, "embedding", 2, &[Some(vec![3.0, 3.0]), None]); + let seg3 = bucket_dir.join("emb-3.vector.parquet"); + write_fixed_size_list_parquet( + &seg3, + "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", "vec_rolled_rr_t"), + table_path, + table_schema, + None, + ); + + // Select absolute rows [1, 3] -> rows at index 1,2,3 (ids 2,3,4; + // embeddings [2,2],[3,3],null). This window straddles the emb-1/emb-2 boundary. + 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, + 6, + 1, + normal_path.metadata().unwrap().len() as i64, + Some(vec!["id"]), + ), + data_file_meta_with_path( + "emb-1.vector.parquet", + 0, + 2, + 1, + seg1.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + data_file_meta_with_path( + "emb-2.vector.parquet", + 2, + 2, + 1, + seg2.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + data_file_meta_with_path( + "emb-3.vector.parquet", + 4, + 2, + 1, + seg3.metadata().unwrap().len() as i64, + Some(vec!["embedding"]), + ), + ]) + .with_row_ranges(vec![RowRange::new(1, 3)]) + .build() + .unwrap(); + + let read = TableRead::new(&table, table.schema().fields().to_vec(), Vec::new()); + let batches = read + .to_arrow(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + + assert_eq!(collect_int_values(&batches, "id"), vec![2, 3, 4]); + assert_fixed_size_list( + &batches, + "embedding", + 2, + &[Some(vec![2.0, 2.0]), Some(vec![3.0, 3.0]), None], + ); + } + /// (6) Row-range mismatch: normal file row_count=3 but `.vector.parquet` row_count=2 /// must surface as DataInvalid. #[tokio::test] @@ -2385,6 +3274,7 @@ mod tests { ResolvedFileInfo { field_ids, data_fields: None, + normalized_write_cols: None, } }