From c81ec9014273d9bf92dd5431b725db7f11f33b11 Mon Sep 17 00:00:00 2001 From: yantian Date: Fri, 14 Aug 2026 09:29:14 +0800 Subject: [PATCH 1/2] fix(vindex): avoid full index read for refine metric --- .../paimon/src/table/vector_search_builder.rs | 103 ++++++++++++++++-- 1 file changed, 95 insertions(+), 8 deletions(-) diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 1295db93..2ae81443 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -17,7 +17,7 @@ use crate::arrow::format::FilePredicates; use crate::arrow::residual::{evaluate_predicates_mask, widen_scan_fields}; -use crate::io::FileIO; +use crate::io::{FileIO, FileRead}; use crate::lumina::reader::LuminaVectorGlobalIndexReader; use crate::lumina::{ is_lumina_index_type, LuminaIndexMeta, LuminaVectorIndexOptions, LuminaVectorMetric, @@ -64,6 +64,7 @@ use crate::vindex::{is_vindex_index_type, VindexVectorIndexOptions}; use arrow_array::{Array, FixedSizeListArray, Float32Array, Int64Array, ListArray, RecordBatch}; use arrow_select::interleave::interleave_record_batch; use futures::{stream, TryStreamExt}; +use paimon_vindex_core::diskann_io::DISKANN_HEADER_SIZE; use paimon_vindex_core::distance::MetricType; use paimon_vindex_core::index::VectorIndexReader as VIndexReader; use paimon_vindex_core::io::SeekRead; @@ -2534,16 +2535,39 @@ async fn resolve_raw_vector_metric( } } VectorIndexBackend::Vindex => { + if let Some(index_meta) = global_meta.index_meta.as_ref() { + if let Ok(options) = + serde_json::from_slice::>(index_meta) + { + if let Some(metric) = options.get("metric") { + return RawVectorMetric::parse(metric); + } + } + } let path = format!("{table_path}/{INDEX_DIR}/{}", entry.index_file.file_name); let input = file_io.new_input(&path)?; - let bytes = input.read().await.map_err(|e| crate::Error::DataInvalid { - message: format!( - "Failed to read vindex index file '{}' for raw search metric: {}", - entry.index_file.file_name, e - ), - source: None, + let file_reader = input + .reader() + .await + .map_err(|e| crate::Error::DataInvalid { + message: format!( + "Failed to read vindex index file '{}' for raw search metric: {}", + entry.index_file.file_name, e + ), + source: Some(Box::new(e)), + })?; + let header_size = + (entry.index_file.file_size as u64).min(DISKANN_HEADER_SIZE as u64); + let bytes = file_reader.read(0..header_size).await.map_err(|e| { + crate::Error::DataInvalid { + message: format!( + "Failed to read vindex index file '{}' for raw search metric: {}", + entry.index_file.file_name, e + ), + source: Some(Box::new(e)), + } })?; - let reader = VIndexReader::open(Cursor::new(bytes.to_vec())).map_err(|e| { + let reader = VIndexReader::open(Cursor::new(bytes)).map_err(|e| { crate::Error::DataInvalid { message: format!( "Failed to open paimon-vindex-core reader for raw search metric: {}", @@ -3167,6 +3191,69 @@ mod tests { ); } + #[tokio::test] + async fn test_resolve_raw_vector_metric_uses_vindex_manifest_metadata() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let mut entry = make_lumina_entry("missing.idx", IVF_FLAT_IDENTIFIER, FileKind::Add, 2); + let index_meta = serde_json::to_vec(&HashMap::from([( + "metric".to_string(), + "cosine".to_string(), + )])) + .unwrap(); + entry + .index_file + .global_index_meta + .as_mut() + .unwrap() + .index_meta = Some(index_meta); + + let metric = resolve_raw_vector_metric( + &file_io, + "memory:///test_table", + &HashMap::new(), + &[entry], + 2, + "embedding", + ) + .await + .unwrap(); + + assert_eq!(metric, RawVectorMetric::Cosine); + } + + #[tokio::test] + async fn test_resolve_raw_vector_metric_reads_vindex_header_for_legacy_metadata() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let index = build_vindex_segment_bytes("inner_product"); + file_io + .new_output("memory:///test_table/index/test.idx") + .unwrap() + .write(bytes::Bytes::from(index.clone())) + .await + .unwrap(); + let mut entry = make_lumina_entry("test.idx", IVF_FLAT_IDENTIFIER, FileKind::Add, 2); + entry.index_file.file_size = index.len() as i64; + entry + .index_file + .global_index_meta + .as_mut() + .unwrap() + .index_meta = Some(b"{}".to_vec()); + + let metric = resolve_raw_vector_metric( + &file_io, + "memory:///test_table", + &HashMap::new(), + &[entry], + 2, + "embedding", + ) + .await + .unwrap(); + + assert_eq!(metric, RawVectorMetric::InnerProduct); + } + #[test] fn test_configured_refine_factor_precedence_and_aliases() { let table_options = HashMap::from([( From d4f7b07e24f9f726b5f9558c87da05dde5637324 Mon Sep 17 00:00:00 2001 From: yantian Date: Fri, 14 Aug 2026 10:27:06 +0800 Subject: [PATCH 2/2] fix(vindex): fall back when refine metric metadata is invalid --- .../paimon/src/table/vector_search_builder.rs | 92 ++++++++++--------- 1 file changed, 50 insertions(+), 42 deletions(-) diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 2ae81443..21c308bf 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -2540,33 +2540,35 @@ async fn resolve_raw_vector_metric( serde_json::from_slice::>(index_meta) { if let Some(metric) = options.get("metric") { - return RawVectorMetric::parse(metric); + if let Some(metric) = + RawVectorMetric::parse_normalized(&normalize_metric(metric)) + { + return Ok(metric); + } } } } let path = format!("{table_path}/{INDEX_DIR}/{}", entry.index_file.file_name); let input = file_io.new_input(&path)?; - let file_reader = input - .reader() - .await - .map_err(|e| crate::Error::DataInvalid { - message: format!( - "Failed to read vindex index file '{}' for raw search metric: {}", - entry.index_file.file_name, e - ), - source: Some(Box::new(e)), - })?; - let header_size = - (entry.index_file.file_size as u64).min(DISKANN_HEADER_SIZE as u64); - let bytes = file_reader.read(0..header_size).await.map_err(|e| { - crate::Error::DataInvalid { - message: format!( - "Failed to read vindex index file '{}' for raw search metric: {}", - entry.index_file.file_name, e - ), - source: Some(Box::new(e)), - } - })?; + let read_error = |e| crate::Error::DataInvalid { + message: format!( + "Failed to read vindex index file '{}' for raw search metric: {}", + entry.index_file.file_name, e + ), + source: Some(Box::new(e)), + }; + let header_size = if entry.index_file.file_size > 0 { + (entry.index_file.file_size as u64).min(DISKANN_HEADER_SIZE as u64) + } else { + input + .metadata() + .await + .map_err(&read_error)? + .size + .min(DISKANN_HEADER_SIZE as u64) + }; + let file_reader = input.reader().await.map_err(&read_error)?; + let bytes = file_reader.read(0..header_size).await.map_err(read_error)?; let reader = VIndexReader::open(Cursor::new(bytes)).map_err(|e| { crate::Error::DataInvalid { message: format!( @@ -3222,7 +3224,7 @@ mod tests { } #[tokio::test] - async fn test_resolve_raw_vector_metric_reads_vindex_header_for_legacy_metadata() { + async fn test_resolve_raw_vector_metric_falls_back_to_vindex_header() { let file_io = FileIOBuilder::new("memory").build().unwrap(); let index = build_vindex_segment_bytes("inner_product"); file_io @@ -3231,27 +3233,33 @@ mod tests { .write(bytes::Bytes::from(index.clone())) .await .unwrap(); - let mut entry = make_lumina_entry("test.idx", IVF_FLAT_IDENTIFIER, FileKind::Add, 2); - entry.index_file.file_size = index.len() as i64; - entry - .index_file - .global_index_meta - .as_mut() - .unwrap() - .index_meta = Some(b"{}".to_vec()); + for (file_size, index_meta) in [ + (index.len() as i64, br#"{"metric":"euclidean"}"#.to_vec()), + (0, b"{}".to_vec()), + (-1, b"{}".to_vec()), + ] { + let mut entry = make_lumina_entry("test.idx", IVF_FLAT_IDENTIFIER, FileKind::Add, 2); + entry.index_file.file_size = file_size; + entry + .index_file + .global_index_meta + .as_mut() + .unwrap() + .index_meta = Some(index_meta); - let metric = resolve_raw_vector_metric( - &file_io, - "memory:///test_table", - &HashMap::new(), - &[entry], - 2, - "embedding", - ) - .await - .unwrap(); + let metric = resolve_raw_vector_metric( + &file_io, + "memory:///test_table", + &HashMap::new(), + &[entry], + 2, + "embedding", + ) + .await + .unwrap(); - assert_eq!(metric, RawVectorMetric::InnerProduct); + assert_eq!(metric, RawVectorMetric::InnerProduct); + } } #[test]