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 d73c1501 feat(table): support the Lumina/DiskANN backend in
primary-key vector read (#555)
d73c1501 is described below
commit d73c1501cc78c2879d8a981f8712650e187022fa
Author: Junrui Lee <[email protected]>
AuthorDate: Mon Jul 20 17:43:42 2026 +0800
feat(table): support the Lumina/DiskANN backend in primary-key vector read
(#555)
---
crates/paimon/src/table/pk_vector_scan.rs | 6 +-
crates/paimon/src/table/vector_search_builder.rs | 181 +++++++++++++++++++----
crates/paimon/src/vindex/pkvector/metric.rs | 29 ++++
3 files changed, 183 insertions(+), 33 deletions(-)
diff --git a/crates/paimon/src/table/pk_vector_scan.rs
b/crates/paimon/src/table/pk_vector_scan.rs
index aebea91b..c4947240 100644
--- a/crates/paimon/src/table/pk_vector_scan.rs
+++ b/crates/paimon/src/table/pk_vector_scan.rs
@@ -330,9 +330,9 @@ fn plan_from_inputs(
source_meta,
path,
file_size,
- // Not consumed on the search path: the vindex reader loads its
- // metadata from the index file bytes and ignores this field,
so an
- // absent value defaulting to an empty vec is acceptable.
+ // The Lumina reader consumes this as its serialized index
+ // metadata; the vindex reader ignores it and loads metadata
from
+ // the segment file bytes. Absent value defaults to an empty
vec.
index_meta: gim.index_meta.clone().unwrap_or_default(),
});
}
diff --git a/crates/paimon/src/table/vector_search_builder.rs
b/crates/paimon/src/table/vector_search_builder.rs
index 67579f1f..992c4cd4 100644
--- a/crates/paimon/src/table/vector_search_builder.rs
+++ b/crates/paimon/src/table/vector_search_builder.rs
@@ -437,13 +437,29 @@ impl<'a> VectorSearchBuilder<'a> {
let search_mode = core.global_index_search_mode()?;
let skip_exact_fallback = search_mode == GlobalIndexSearchMode::Fast;
- let plan = PkVectorScan::new(self.table, field_id, index_type,
self.filter.clone())
- .plan()
- .await?;
+ let plan = PkVectorScan::new(
+ self.table,
+ field_id,
+ index_type.clone(),
+ self.filter.clone(),
+ )
+ .plan()
+ .await?;
if plan.splits.is_empty() {
return Ok((Vec::new(), plan, metric));
}
+ // Resolve the vector index backend from the single configured index
type.
+ // Java enforces one index type per PK table and Rust filters segments
to
+ // it, so one backend serves every segment. Computed after the
empty-plan
+ // return so an empty table never errors on an unrecognized type.
+ let backend =
VectorIndexBackend::from_index_type(&index_type).ok_or_else(|| {
+ crate::Error::DataInvalid {
+ message: format!("unsupported PK vector index backend/type:
'{index_type}'"),
+ source: None,
+ }
+ })?;
+
// Production data-file reader, mirroring
`table_read.rs::new_data_file_reader`
// but projecting only the vector column with no predicates.
let reader = DataFileReader::new(
@@ -460,7 +476,7 @@ impl<'a> VectorSearchBuilder<'a> {
let segment_bytes = preload_segment_bytes(self.table.file_io(),
&plan.splits).await?;
// Fail loud on a config/segment metric mismatch before scoring,
mirroring
// Java `PkVectorAnnSegmentSearcher.search`.
- verify_pk_vector_segment_metrics(&plan.splits, &segment_bytes,
metric)?;
+ verify_pk_vector_segment_metrics(&plan.splits, &segment_bytes, metric,
backend)?;
let options = {
let mut o = self.table.schema().options().clone();
o.extend(self.options.clone());
@@ -482,8 +498,18 @@ impl<'a> VectorSearchBuilder<'a> {
segment.file_size,
segment.index_meta.clone(),
);
- let mut reader = VindexVectorGlobalIndexReader::new(io_meta,
options.clone());
- reader.visit_vector_search(search, |_| Ok(Cursor::new(data)))
+ match backend {
+ VectorIndexBackend::Lumina => {
+ let mut reader =
+ LuminaVectorGlobalIndexReader::new(io_meta,
options.clone());
+ reader.visit_vector_search(search, |_|
Ok(Cursor::new(data)))
+ }
+ VectorIndexBackend::Vindex => {
+ let mut reader =
+ VindexVectorGlobalIndexReader::new(io_meta,
options.clone());
+ reader.visit_vector_search(search, |_|
Ok(Cursor::new(data)))
+ }
+ }
});
let ann_searcher = VindexAnnSearcher::new(field_name, scorer);
@@ -1200,6 +1226,7 @@ fn verify_pk_vector_segment_metrics(
splits: &[PkVectorSearchSplit],
segment_bytes: &HashMap<String, Vec<u8>>,
configured: VectorSearchMetric,
+ backend: VectorIndexBackend,
) -> crate::Result<()> {
let mut checked: HashSet<&str> = HashSet::new();
for split in splits {
@@ -1207,27 +1234,37 @@ fn verify_pk_vector_segment_metrics(
if !checked.insert(segment.path.as_str()) {
continue;
}
- let bytes =
- segment_bytes
- .get(&segment.path)
- .ok_or_else(|| crate::Error::DataInvalid {
- message: format!(
- "missing preloaded ANN bytes for segment '{}'",
- segment.path
- ),
- source: None,
+ let segment_metric = match backend {
+ VectorIndexBackend::Lumina => {
+ // Lumina records its metric in the serialized index
metadata
+ // (`index_meta`), not in the segment file bytes.
+ let lumina_metric =
+
LuminaIndexMeta::deserialize(&segment.index_meta)?.metric()?;
+ VectorSearchMetric::from_lumina(lumina_metric)
+ }
+ VectorIndexBackend::Vindex => {
+ let bytes = segment_bytes.get(&segment.path).ok_or_else(||
{
+ crate::Error::DataInvalid {
+ message: format!(
+ "missing preloaded ANN bytes for segment '{}'",
+ segment.path
+ ),
+ source: None,
+ }
})?;
- let reader =
VIndexReader::open(Cursor::new(bytes.clone())).map_err(|e| {
- crate::Error::DataInvalid {
- message: format!(
- "failed to open ANN index file '{}' for metric check:
{e}",
- segment.path
- ),
- source: Some(Box::new(e)),
+ let reader =
VIndexReader::open(Cursor::new(bytes.clone())).map_err(|e| {
+ crate::Error::DataInvalid {
+ message: format!(
+ "failed to open ANN index file '{}' for metric
check: {e}",
+ segment.path
+ ),
+ source: Some(Box::new(e)),
+ }
+ })?;
+ VectorSearchMetric::from_vindex(reader.metadata().metric)
}
- })?;
- let segment_metric = reader.metadata().metric;
- if VectorSearchMetric::from_vindex(segment_metric) != configured {
+ };
+ if segment_metric != configured {
return Err(crate::Error::DataInvalid {
message: format!(
"ANN segment metric {} does not match configured
metric {}",
@@ -2971,14 +3008,94 @@ mod tests {
split
}
+ fn pk_split_with_lumina_segment(path: &str, metric: &str) ->
PkVectorSearchSplit {
+ let mut split = pk_search_split(0, vec![pk_data_file("file-a", 3,
Some(0))]);
+ let source_meta = crate::spec::PkVectorSourceMeta::new(
+ 1,
+ vec![crate::spec::PkVectorSourceFile::new("file-a".to_string(),
3).unwrap()],
+ )
+ .unwrap();
+ let mut segment = BucketAnnSegment::for_test(source_meta);
+ segment.path = path.to_string();
+ // Lumina stores its metric in the serialized index metadata blob, not
in
+ // the segment file bytes. `deserialize` requires both keys present.
+ let meta = crate::lumina::LuminaIndexMeta::new(HashMap::from([
+ ("index.dimension".to_string(), "2".to_string()),
+ ("distance.metric".to_string(), metric.to_string()),
+ ]));
+ segment.index_meta = meta.serialize().unwrap();
+ split.ann_segments = vec![segment];
+ split
+ }
+
+ #[test]
+ fn verify_pk_vector_segment_metrics_accepts_matching_lumina_metric() {
+ // Lumina segment metadata says cosine; configured cosine => Ok. No
segment
+ // file bytes are needed on the Lumina path.
+ let splits = vec![pk_split_with_lumina_segment("seg-lumina",
"cosine")];
+ let segment_bytes = HashMap::new();
+ verify_pk_vector_segment_metrics(
+ &splits,
+ &segment_bytes,
+ VectorSearchMetric::Cosine,
+ VectorIndexBackend::Lumina,
+ )
+ .expect("matching lumina metric must pass");
+ }
+
+ #[test]
+ fn verify_pk_vector_segment_metrics_rejects_mismatched_lumina_metric() {
+ // Lumina segment metadata says l2; configured inner_product => fail
loud,
+ // naming both metrics.
+ let splits = vec![pk_split_with_lumina_segment("seg-lumina", "l2")];
+ let segment_bytes = HashMap::new();
+ let err = verify_pk_vector_segment_metrics(
+ &splits,
+ &segment_bytes,
+ VectorSearchMetric::InnerProduct,
+ VectorIndexBackend::Lumina,
+ )
+ .expect_err("mismatched lumina metric must fail loud");
+ assert!(
+ matches!(err, crate::Error::DataInvalid { ref message, .. }
+ if message.contains("does not match configured metric")
+ && message.contains("l2")
+ && message.contains("inner_product")),
+ "unexpected error: {err:?}"
+ );
+ }
+
+ #[test]
+ fn from_index_type_classifies_lumina_and_vindex() {
+ assert_eq!(
+ VectorIndexBackend::from_index_type("lumina"),
+ Some(VectorIndexBackend::Lumina)
+ );
+ assert_eq!(
+ VectorIndexBackend::from_index_type("lumina-vector-ann"),
+ Some(VectorIndexBackend::Lumina)
+ );
+ assert_eq!(
+ VectorIndexBackend::from_index_type("ivf-flat"),
+ Some(VectorIndexBackend::Vindex)
+ );
+ // `diskann` is Lumina's internal index type, not a top-level index
type.
+ assert_eq!(VectorIndexBackend::from_index_type("diskann"), None);
+ }
+
#[test]
fn verify_pk_vector_segment_metrics_accepts_matching_metric() {
// Real IVF segment trained with L2; configured metric L2 => Ok.
let bytes = build_vindex_segment_bytes("l2");
let splits = vec![pk_split_with_segment("seg-l2")];
let segment_bytes = HashMap::from([("seg-l2".to_string(), bytes)]);
- verify_pk_vector_segment_metrics(&splits, &segment_bytes,
VectorSearchMetric::L2)
- .expect("matching metric must pass");
+ verify_pk_vector_segment_metrics(
+ &splits,
+ &segment_bytes,
+ VectorSearchMetric::L2,
+ VectorIndexBackend::Vindex,
+ )
+ .expect("matching metric must pass");
}
#[test]
@@ -2987,9 +3104,13 @@ mod tests {
let bytes = build_vindex_segment_bytes("l2");
let splits = vec![pk_split_with_segment("seg-l2")];
let segment_bytes = HashMap::from([("seg-l2".to_string(), bytes)]);
- let err =
- verify_pk_vector_segment_metrics(&splits, &segment_bytes,
VectorSearchMetric::Cosine)
- .expect_err("mismatched metric must fail loud");
+ let err = verify_pk_vector_segment_metrics(
+ &splits,
+ &segment_bytes,
+ VectorSearchMetric::Cosine,
+ VectorIndexBackend::Vindex,
+ )
+ .expect_err("mismatched metric must fail loud");
assert!(
matches!(err, crate::Error::DataInvalid { ref message, .. }
if message.contains("does not match configured metric")
diff --git a/crates/paimon/src/vindex/pkvector/metric.rs
b/crates/paimon/src/vindex/pkvector/metric.rs
index 7d654e58..59aa24c5 100644
--- a/crates/paimon/src/vindex/pkvector/metric.rs
+++ b/crates/paimon/src/vindex/pkvector/metric.rs
@@ -16,6 +16,7 @@
// under the License.
use super::data_invalid;
+use crate::lumina::LuminaVectorMetric;
use std::cmp::Ordering;
/// Order two distances the way Java `Float.compare` does: every NaN sorts
after
@@ -62,6 +63,17 @@ impl VectorSearchMetric {
}
}
+ /// Map a Lumina metric to this enum, symmetric to `from_vindex`. Lets the
+ /// read path compare the metric a Lumina segment was built with against
the
+ /// configured metric.
+ pub(crate) fn from_lumina(metric: LuminaVectorMetric) -> Self {
+ match metric {
+ LuminaVectorMetric::L2 => Self::L2,
+ LuminaVectorMetric::Cosine => Self::Cosine,
+ LuminaVectorMetric::InnerProduct => Self::InnerProduct,
+ }
+ }
+
/// Normalize, validate, and map to the enum. Errors on an unsupported
metric.
pub(crate) fn parse(metric: &str) -> crate::Result<Self> {
match normalize_metric(metric).as_str() {
@@ -320,4 +332,21 @@ mod tests {
VectorSearchMetric::InnerProduct
);
}
+
+ #[test]
+ fn test_from_lumina_maps_every_variant() {
+ use crate::lumina::LuminaVectorMetric;
+ assert_eq!(
+ VectorSearchMetric::from_lumina(LuminaVectorMetric::L2),
+ VectorSearchMetric::L2
+ );
+ assert_eq!(
+ VectorSearchMetric::from_lumina(LuminaVectorMetric::Cosine),
+ VectorSearchMetric::Cosine
+ );
+ assert_eq!(
+ VectorSearchMetric::from_lumina(LuminaVectorMetric::InnerProduct),
+ VectorSearchMetric::InnerProduct
+ );
+ }
}