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 b4377f3 [table] Support reading dedicated and rolling (multi-segment)
vector files (#423)
b4377f3 is described below
commit b4377f3eabbe7be0f389365e162ba0108369b229
Author: Junrui Lee <[email protected]>
AuthorDate: Tue Jun 30 16:44:29 2026 +0800
[table] Support reading dedicated and rolling (multi-segment) vector files
(#423)
Route dedicated `*.vector.<format>` 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
---
crates/paimon/src/table/data_evolution_reader.rs | 1959 ++++++++++++++++++++--
1 file changed, 1853 insertions(+), 106 deletions(-)
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index 487ed95..8f4a5cb 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.<format>`).
+/// 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 VectorBunch source. Check
this
+ // BEFORE the single-file early-return.
+ if files
+ .iter()
+ .any(|file| is_vector_store_file_name(&file.file_name))
+ {
+ return false;
+ }
if files.len() <= 1 {
return true;
}
@@ -509,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<DataFileMeta>,
+ data_fields: Option<Vec<DataField>>,
+ row_ranges: Option<Vec<RowRange>>,
+) -> crate::Result<ArrowRecordBatchStream> {
+ 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)]
@@ -552,11 +588,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,
});
}
@@ -591,6 +629,7 @@ impl PreparedMergeGroup {
struct ResolvedFileInfo {
field_ids: Vec<i32>,
data_fields: Option<Vec<DataField>>,
+ normalized_write_cols: Option<Vec<String>>,
}
async fn load_file_infos(
@@ -602,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)
@@ -642,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<Vec<String>> {
+ 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::<crate::Result<_>>()?;
+
+ indexed.sort_by_key(|(pos, _)| *pos);
+ Ok(indexed.into_iter().map(|(_, name)| name).collect())
+}
+
#[derive(Debug, Clone)]
struct SourcePlan {
sources: Vec<FieldSource>,
@@ -655,7 +753,9 @@ fn build_source_plan(
blob_descriptor_fields: &HashSet<String>,
) -> crate::Result<SourcePlan> {
let mut sources = Vec::new();
- let mut normal_source_indices: HashMap<usize, usize> = HashMap::new();
+ let mut normal_providers: HashMap<i32, usize> = HashMap::new(); //
field_id -> source_idx
+ let mut vector_field_providers: HashMap<i32, usize> = HashMap::new(); //
field_id -> source_idx
+ let mut vector_bunch_indices: HashMap<(i64, String, Vec<String>), usize> =
HashMap::new();
let mut blob_source_indices: HashMap<i32, usize> = HashMap::new();
let mut expected_blob_row_count: Option<i64> = None;
@@ -688,6 +788,65 @@ 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). 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 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 {
expected_blob_row_count = Some(file.row_count);
let source_idx = sources.len();
@@ -696,7 +855,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 +868,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. bunch; fall back to a normal data
file
+ // (PR 2 inline-vector compatibility path).
+ vector_field_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 {
@@ -741,31 +906,30 @@ 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,
})
}
-fn select_normal_provider(
- files: &[DataFileMeta],
- file_infos: &[ResolvedFileInfo],
- normal_source_indices: &HashMap<usize, usize>,
- field_id: i32,
-) -> Option<usize> {
- 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<i32> {
if info.field_ids.len() != 1 {
return Err(Error::DataInvalid {
@@ -788,6 +952,11 @@ enum FieldSource {
data_fields: Option<Vec<DataField>>,
read_fields: Vec<DataField>,
},
+ VectorBunch {
+ bunch: VectorBunch,
+ data_fields: Option<Vec<DataField>>,
+ read_fields: Vec<DataField>,
+ },
BlobBunch {
bunch: BlobBunch,
data_fields: Option<Vec<DataField>>,
@@ -799,6 +968,7 @@ impl FieldSource {
fn read_fields(&self) -> &[DataField] {
match self {
FieldSource::DataFile { read_fields, .. }
+ | FieldSource::VectorBunch { read_fields, .. }
| FieldSource::BlobBunch { read_fields, .. } => read_fields,
}
}
@@ -806,6 +976,7 @@ impl FieldSource {
fn add_read_field(&mut self, field: DataField) -> usize {
let read_fields = match self {
FieldSource::DataFile { read_fields, .. }
+ | FieldSource::VectorBunch { read_fields, .. }
| FieldSource::BlobBunch { read_fields, .. } => read_fields,
};
if let Some(offset) = read_fields
@@ -822,7 +993,14 @@ impl FieldSource {
fn blob_bunch_mut(&mut self) -> Option<&mut BlobBunch> {
match self {
FieldSource::BlobBunch { bunch, .. } => Some(bunch),
- FieldSource::DataFile { .. } => 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,
}
}
}
@@ -936,45 +1114,196 @@ impl BlobBunch {
}
}
+/// Aggregates rolled `.vector.<format>` 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<DataFileMeta>,
+ schema_id: i64,
+ format_suffix: String,
+ normalized_write_cols: Vec<String>,
+ 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<String>,
+ ) -> 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<DataFileMeta>) ->
crate::Result<Vec<DataFileMeta>> {
- 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
{
+ normal_files.sort_by_key(|f| std::cmp::Reverse(f.max_sequence_number));
+
+ // Vector files: sort by first_row_id asc, then max_sequence_number desc
(like blobs).
+ // They are NOT validated against the normal-file row range — rolled
segments are
+ // slices with their own ranges. They DO require first_row_id.
+ if vector_files.iter().any(|file| file.first_row_id.is_none()) {
+ return Err(Error::DataInvalid {
+ message: "All vector files in a field merge split should have
first_row_id".to_string(),
+ source: None,
+ });
+ }
+ vector_files.sort_by(|left, right| {
+ 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))
+ });
+
+ // Normal files share the anchor's row range. Validate normal files ONLY
(vectors removed).
+ let mut range_ref: Option<(i64, i64)> = None;
+ for file in normal_files.iter() {
+ let first_row_id = file.first_row_id.ok_or_else(|| Error::DataInvalid {
message: "All data 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 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 +1311,10 @@ fn normalize_merge_group(files: Vec<DataFileMeta>) ->
crate::Result<Vec<DataFile
});
}
- data_files.extend(blob_files);
- Ok(data_files)
+ let mut out = normal_files;
+ out.extend(vector_files);
+ out.extend(blob_files);
+ Ok(out)
}
fn count_selected_rows(
@@ -1006,9 +1337,11 @@ mod tests {
use crate::catalog::Identifier;
use crate::io::FileIOBuilder;
use crate::spec::stats::BinaryTableStats;
- use crate::spec::{BinaryRow, BlobType, IntType, Schema, TableSchema};
+ use crate::spec::{BinaryRow, BlobType, FloatType, IntType, Schema,
TableSchema, VectorType};
use crate::table::{DataSplitBuilder, Table, TableRead};
- use arrow_array::{Array, BinaryArray, Int32Array, RecordBatch};
+ use arrow_array::{
+ Array, BinaryArray, FixedSizeListArray, Float32Array, Int32Array,
RecordBatch,
+ };
use futures::TryStreamExt;
use std::fs;
use std::path::{Path, PathBuf};
@@ -1030,35 +1363,245 @@ mod tests {
use test_utils::{local_file_path, write_int_parquet_file};
#[test]
- fn test_normalize_merge_group_orders_blob_files_after_data_files() {
+ 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("file1.parquet", 1, 10, 1, None),
- data_file("file2.blob", 1, 1, 1, Some(vec!["payload"])),
- data_file("file3.blob", 1, 1, 3, Some(vec!["payload"])),
- data_file("file4.blob", 2, 9, 1, Some(vec!["payload"])),
- data_file("file7.parquet", 1, 10, 3, None),
+ 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: 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();
- let normalized = normalize_merge_group(files).unwrap();
- let file_names: Vec<&str> = normalized
- .iter()
- .map(|file| file.file_name.as_str())
- .collect();
- assert_eq!(
- file_names,
- vec![
- "file7.parquet",
- "file1.parquet",
- "file3.blob",
- "file2.blob",
- "file4.blob",
- ]
- );
+ // 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_blob_bunch_ignores_same_first_row_id_with_lower_sequence() {
- let mut bunch = BlobBunch::new(1000);
+ 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, "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 { .. })));
+ }
+
+ #[test]
+ fn test_normalize_merge_group_orders_blob_files_after_data_files() {
+ let files = vec![
+ data_file("file1.parquet", 1, 10, 1, None),
+ data_file("file2.blob", 1, 1, 1, Some(vec!["payload"])),
+ data_file("file3.blob", 1, 1, 3, Some(vec!["payload"])),
+ data_file("file4.blob", 2, 9, 1, Some(vec!["payload"])),
+ data_file("file7.parquet", 1, 10, 3, None),
+ ];
+
+ let normalized = normalize_merge_group(files).unwrap();
+ let file_names: Vec<&str> = normalized
+ .iter()
+ .map(|file| file.file_name.as_str())
+ .collect();
+ assert_eq!(
+ file_names,
+ vec![
+ "file7.parquet",
+ "file1.parquet",
+ "file3.blob",
+ "file2.blob",
+ "file4.blob",
+ ]
+ );
+ }
+
+ #[test]
+ fn test_normalize_merge_group_orders_vector_files_between_data_and_blob() {
+ // Discriminating fixture: the vector file has a HIGHER
max_sequence_number than
+ // the normal file and is listed first. Old two-group code sorted it
among the
+ // "data files" by Reverse(seq), yielding [v1, d1, ...]; the three-way
split must
+ // force normal -> 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_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"])),
+ vector_no_rid,
+ ];
+ 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);
bunch
.add(data_file(
"blob-high.blob",
@@ -1077,6 +1620,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);
@@ -1179,6 +1747,204 @@ mod tests {
);
}
+ #[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("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("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_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("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("v-high.vector.parquet", 0, 10, 3,
Some(vec!["emb"])),
+ &["emb".to_string()],
+ )
+ .unwrap_err();
+ assert!(matches!(err, Error::DataInvalid { .. }));
+ }
+
+ #[test]
+ 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("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);
+ }
+
+ #[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_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("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("v1.parquet", 0, 10, 1, Some(vec!["emb"])),
+ &["emb".to_string()],
+ )
+ .unwrap_err();
+ assert!(matches!(err, Error::DataInvalid { .. }));
+ }
+
+ #[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]
fn test_build_source_plan_picks_latest_blob_segments() {
let files = vec![
@@ -1228,7 +1994,9 @@ mod tests {
vec!["blob5.blob", "blob9.blob", "blob7.blob",
"blob8.blob"]
);
}
- FieldSource::DataFile { .. } => panic!("expected blob bunch
source"),
+ FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => {
+ panic!("expected blob bunch source")
+ }
}
}
@@ -1308,7 +2076,9 @@ mod tests {
vec!["blob5.blob", "blob9.blob", "blob7.blob",
"blob8.blob"]
);
}
- FieldSource::DataFile { .. } => panic!("expected blob bunch
source"),
+ FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => {
+ panic!("expected blob bunch source")
+ }
}
match &source_plan.sources[2] {
@@ -1323,7 +2093,9 @@ mod tests {
vec!["blob15.blob", "blob19.blob", "blob17.blob",
"blob18.blob"]
);
}
- FieldSource::DataFile { .. } => panic!("expected blob bunch
source"),
+ FieldSource::DataFile { .. } | FieldSource::VectorBunch { .. } => {
+ panic!("expected blob bunch source")
+ }
}
}
@@ -1524,10 +2296,985 @@ mod tests {
);
}
+ fn write_fixed_size_list_parquet(
+ path: &std::path::Path,
+ col: &str,
+ dim: i32,
+ rows: &[Option<Vec<f32>>],
+ ) {
+ 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<FLOAT, dim>` 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())
+ }
+
+ #[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(
+ batches: &[RecordBatch],
+ column_name: &str,
+ expected_dim: i32,
+ expected: &[Option<Vec<f32>>],
+ ) {
+ 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::<FixedSizeListArray>()
+ .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::<Float32Array>().unwrap();
+ let got: Vec<f32> = (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::<Vec<_>>()
+ .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, 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::<Vec<_>>()
+ .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 VectorBunch 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::<Vec<_>>()
+ .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::<Vec<_>>()
+ .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::<Vec<_>>()
+ .await
+ .unwrap();
+
+ 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::<Vec<_>>()
+ .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::<Vec<_>>()
+ .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]
+ 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::<Vec<_>>().await,
+ Err(e) => Err(e),
+ };
+ assert!(
+ matches!(collected, Err(Error::DataInvalid { .. })),
+ "expected DataInvalid, got {collected:?}"
+ );
+ }
+
fn resolved_info(field_ids: Vec<i32>) -> ResolvedFileInfo {
ResolvedFileInfo {
field_ids,
data_fields: None,
+ normalized_write_cols: None,
}
}