Skip to content

Commit d73c150

Browse files
authored
feat(table): support the Lumina/DiskANN backend in primary-key vector read (#555)
1 parent ddd8792 commit d73c150

3 files changed

Lines changed: 183 additions & 33 deletions

File tree

crates/paimon/src/table/pk_vector_scan.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -330,9 +330,9 @@ fn plan_from_inputs(
330330
source_meta,
331331
path,
332332
file_size,
333-
// Not consumed on the search path: the vindex reader loads its
334-
// metadata from the index file bytes and ignores this field, so an
335-
// absent value defaulting to an empty vec is acceptable.
333+
// The Lumina reader consumes this as its serialized index
334+
// metadata; the vindex reader ignores it and loads metadata from
335+
// the segment file bytes. Absent value defaults to an empty vec.
336336
index_meta: gim.index_meta.clone().unwrap_or_default(),
337337
});
338338
}

crates/paimon/src/table/vector_search_builder.rs

Lines changed: 151 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -437,13 +437,29 @@ impl<'a> VectorSearchBuilder<'a> {
437437
let search_mode = core.global_index_search_mode()?;
438438
let skip_exact_fallback = search_mode == GlobalIndexSearchMode::Fast;
439439

440-
let plan = PkVectorScan::new(self.table, field_id, index_type, self.filter.clone())
441-
.plan()
442-
.await?;
440+
let plan = PkVectorScan::new(
441+
self.table,
442+
field_id,
443+
index_type.clone(),
444+
self.filter.clone(),
445+
)
446+
.plan()
447+
.await?;
443448
if plan.splits.is_empty() {
444449
return Ok((Vec::new(), plan, metric));
445450
}
446451

452+
// Resolve the vector index backend from the single configured index type.
453+
// Java enforces one index type per PK table and Rust filters segments to
454+
// it, so one backend serves every segment. Computed after the empty-plan
455+
// return so an empty table never errors on an unrecognized type.
456+
let backend = VectorIndexBackend::from_index_type(&index_type).ok_or_else(|| {
457+
crate::Error::DataInvalid {
458+
message: format!("unsupported PK vector index backend/type: '{index_type}'"),
459+
source: None,
460+
}
461+
})?;
462+
447463
// Production data-file reader, mirroring `table_read.rs::new_data_file_reader`
448464
// but projecting only the vector column with no predicates.
449465
let reader = DataFileReader::new(
@@ -460,7 +476,7 @@ impl<'a> VectorSearchBuilder<'a> {
460476
let segment_bytes = preload_segment_bytes(self.table.file_io(), &plan.splits).await?;
461477
// Fail loud on a config/segment metric mismatch before scoring, mirroring
462478
// Java `PkVectorAnnSegmentSearcher.search`.
463-
verify_pk_vector_segment_metrics(&plan.splits, &segment_bytes, metric)?;
479+
verify_pk_vector_segment_metrics(&plan.splits, &segment_bytes, metric, backend)?;
464480
let options = {
465481
let mut o = self.table.schema().options().clone();
466482
o.extend(self.options.clone());
@@ -482,8 +498,18 @@ impl<'a> VectorSearchBuilder<'a> {
482498
segment.file_size,
483499
segment.index_meta.clone(),
484500
);
485-
let mut reader = VindexVectorGlobalIndexReader::new(io_meta, options.clone());
486-
reader.visit_vector_search(search, |_| Ok(Cursor::new(data)))
501+
match backend {
502+
VectorIndexBackend::Lumina => {
503+
let mut reader =
504+
LuminaVectorGlobalIndexReader::new(io_meta, options.clone());
505+
reader.visit_vector_search(search, |_| Ok(Cursor::new(data)))
506+
}
507+
VectorIndexBackend::Vindex => {
508+
let mut reader =
509+
VindexVectorGlobalIndexReader::new(io_meta, options.clone());
510+
reader.visit_vector_search(search, |_| Ok(Cursor::new(data)))
511+
}
512+
}
487513
});
488514
let ann_searcher = VindexAnnSearcher::new(field_name, scorer);
489515

@@ -1200,34 +1226,45 @@ fn verify_pk_vector_segment_metrics(
12001226
splits: &[PkVectorSearchSplit],
12011227
segment_bytes: &HashMap<String, Vec<u8>>,
12021228
configured: VectorSearchMetric,
1229+
backend: VectorIndexBackend,
12031230
) -> crate::Result<()> {
12041231
let mut checked: HashSet<&str> = HashSet::new();
12051232
for split in splits {
12061233
for segment in &split.ann_segments {
12071234
if !checked.insert(segment.path.as_str()) {
12081235
continue;
12091236
}
1210-
let bytes =
1211-
segment_bytes
1212-
.get(&segment.path)
1213-
.ok_or_else(|| crate::Error::DataInvalid {
1214-
message: format!(
1215-
"missing preloaded ANN bytes for segment '{}'",
1216-
segment.path
1217-
),
1218-
source: None,
1237+
let segment_metric = match backend {
1238+
VectorIndexBackend::Lumina => {
1239+
// Lumina records its metric in the serialized index metadata
1240+
// (`index_meta`), not in the segment file bytes.
1241+
let lumina_metric =
1242+
LuminaIndexMeta::deserialize(&segment.index_meta)?.metric()?;
1243+
VectorSearchMetric::from_lumina(lumina_metric)
1244+
}
1245+
VectorIndexBackend::Vindex => {
1246+
let bytes = segment_bytes.get(&segment.path).ok_or_else(|| {
1247+
crate::Error::DataInvalid {
1248+
message: format!(
1249+
"missing preloaded ANN bytes for segment '{}'",
1250+
segment.path
1251+
),
1252+
source: None,
1253+
}
12191254
})?;
1220-
let reader = VIndexReader::open(Cursor::new(bytes.clone())).map_err(|e| {
1221-
crate::Error::DataInvalid {
1222-
message: format!(
1223-
"failed to open ANN index file '{}' for metric check: {e}",
1224-
segment.path
1225-
),
1226-
source: Some(Box::new(e)),
1255+
let reader = VIndexReader::open(Cursor::new(bytes.clone())).map_err(|e| {
1256+
crate::Error::DataInvalid {
1257+
message: format!(
1258+
"failed to open ANN index file '{}' for metric check: {e}",
1259+
segment.path
1260+
),
1261+
source: Some(Box::new(e)),
1262+
}
1263+
})?;
1264+
VectorSearchMetric::from_vindex(reader.metadata().metric)
12271265
}
1228-
})?;
1229-
let segment_metric = reader.metadata().metric;
1230-
if VectorSearchMetric::from_vindex(segment_metric) != configured {
1266+
};
1267+
if segment_metric != configured {
12311268
return Err(crate::Error::DataInvalid {
12321269
message: format!(
12331270
"ANN segment metric {} does not match configured metric {}",
@@ -2971,14 +3008,94 @@ mod tests {
29713008
split
29723009
}
29733010

3011+
fn pk_split_with_lumina_segment(path: &str, metric: &str) -> PkVectorSearchSplit {
3012+
let mut split = pk_search_split(0, vec![pk_data_file("file-a", 3, Some(0))]);
3013+
let source_meta = crate::spec::PkVectorSourceMeta::new(
3014+
1,
3015+
vec![crate::spec::PkVectorSourceFile::new("file-a".to_string(), 3).unwrap()],
3016+
)
3017+
.unwrap();
3018+
let mut segment = BucketAnnSegment::for_test(source_meta);
3019+
segment.path = path.to_string();
3020+
// Lumina stores its metric in the serialized index metadata blob, not in
3021+
// the segment file bytes. `deserialize` requires both keys present.
3022+
let meta = crate::lumina::LuminaIndexMeta::new(HashMap::from([
3023+
("index.dimension".to_string(), "2".to_string()),
3024+
("distance.metric".to_string(), metric.to_string()),
3025+
]));
3026+
segment.index_meta = meta.serialize().unwrap();
3027+
split.ann_segments = vec![segment];
3028+
split
3029+
}
3030+
3031+
#[test]
3032+
fn verify_pk_vector_segment_metrics_accepts_matching_lumina_metric() {
3033+
// Lumina segment metadata says cosine; configured cosine => Ok. No segment
3034+
// file bytes are needed on the Lumina path.
3035+
let splits = vec![pk_split_with_lumina_segment("seg-lumina", "cosine")];
3036+
let segment_bytes = HashMap::new();
3037+
verify_pk_vector_segment_metrics(
3038+
&splits,
3039+
&segment_bytes,
3040+
VectorSearchMetric::Cosine,
3041+
VectorIndexBackend::Lumina,
3042+
)
3043+
.expect("matching lumina metric must pass");
3044+
}
3045+
3046+
#[test]
3047+
fn verify_pk_vector_segment_metrics_rejects_mismatched_lumina_metric() {
3048+
// Lumina segment metadata says l2; configured inner_product => fail loud,
3049+
// naming both metrics.
3050+
let splits = vec![pk_split_with_lumina_segment("seg-lumina", "l2")];
3051+
let segment_bytes = HashMap::new();
3052+
let err = verify_pk_vector_segment_metrics(
3053+
&splits,
3054+
&segment_bytes,
3055+
VectorSearchMetric::InnerProduct,
3056+
VectorIndexBackend::Lumina,
3057+
)
3058+
.expect_err("mismatched lumina metric must fail loud");
3059+
assert!(
3060+
matches!(err, crate::Error::DataInvalid { ref message, .. }
3061+
if message.contains("does not match configured metric")
3062+
&& message.contains("l2")
3063+
&& message.contains("inner_product")),
3064+
"unexpected error: {err:?}"
3065+
);
3066+
}
3067+
3068+
#[test]
3069+
fn from_index_type_classifies_lumina_and_vindex() {
3070+
assert_eq!(
3071+
VectorIndexBackend::from_index_type("lumina"),
3072+
Some(VectorIndexBackend::Lumina)
3073+
);
3074+
assert_eq!(
3075+
VectorIndexBackend::from_index_type("lumina-vector-ann"),
3076+
Some(VectorIndexBackend::Lumina)
3077+
);
3078+
assert_eq!(
3079+
VectorIndexBackend::from_index_type("ivf-flat"),
3080+
Some(VectorIndexBackend::Vindex)
3081+
);
3082+
// `diskann` is Lumina's internal index type, not a top-level index type.
3083+
assert_eq!(VectorIndexBackend::from_index_type("diskann"), None);
3084+
}
3085+
29743086
#[test]
29753087
fn verify_pk_vector_segment_metrics_accepts_matching_metric() {
29763088
// Real IVF segment trained with L2; configured metric L2 => Ok.
29773089
let bytes = build_vindex_segment_bytes("l2");
29783090
let splits = vec![pk_split_with_segment("seg-l2")];
29793091
let segment_bytes = HashMap::from([("seg-l2".to_string(), bytes)]);
2980-
verify_pk_vector_segment_metrics(&splits, &segment_bytes, VectorSearchMetric::L2)
2981-
.expect("matching metric must pass");
3092+
verify_pk_vector_segment_metrics(
3093+
&splits,
3094+
&segment_bytes,
3095+
VectorSearchMetric::L2,
3096+
VectorIndexBackend::Vindex,
3097+
)
3098+
.expect("matching metric must pass");
29823099
}
29833100

29843101
#[test]
@@ -2987,9 +3104,13 @@ mod tests {
29873104
let bytes = build_vindex_segment_bytes("l2");
29883105
let splits = vec![pk_split_with_segment("seg-l2")];
29893106
let segment_bytes = HashMap::from([("seg-l2".to_string(), bytes)]);
2990-
let err =
2991-
verify_pk_vector_segment_metrics(&splits, &segment_bytes, VectorSearchMetric::Cosine)
2992-
.expect_err("mismatched metric must fail loud");
3107+
let err = verify_pk_vector_segment_metrics(
3108+
&splits,
3109+
&segment_bytes,
3110+
VectorSearchMetric::Cosine,
3111+
VectorIndexBackend::Vindex,
3112+
)
3113+
.expect_err("mismatched metric must fail loud");
29933114
assert!(
29943115
matches!(err, crate::Error::DataInvalid { ref message, .. }
29953116
if message.contains("does not match configured metric")

crates/paimon/src/vindex/pkvector/metric.rs

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
// under the License.
1717

1818
use super::data_invalid;
19+
use crate::lumina::LuminaVectorMetric;
1920
use std::cmp::Ordering;
2021

2122
/// Order two distances the way Java `Float.compare` does: every NaN sorts after
@@ -62,6 +63,17 @@ impl VectorSearchMetric {
6263
}
6364
}
6465

66+
/// Map a Lumina metric to this enum, symmetric to `from_vindex`. Lets the
67+
/// read path compare the metric a Lumina segment was built with against the
68+
/// configured metric.
69+
pub(crate) fn from_lumina(metric: LuminaVectorMetric) -> Self {
70+
match metric {
71+
LuminaVectorMetric::L2 => Self::L2,
72+
LuminaVectorMetric::Cosine => Self::Cosine,
73+
LuminaVectorMetric::InnerProduct => Self::InnerProduct,
74+
}
75+
}
76+
6577
/// Normalize, validate, and map to the enum. Errors on an unsupported metric.
6678
pub(crate) fn parse(metric: &str) -> crate::Result<Self> {
6779
match normalize_metric(metric).as_str() {
@@ -320,4 +332,21 @@ mod tests {
320332
VectorSearchMetric::InnerProduct
321333
);
322334
}
335+
336+
#[test]
337+
fn test_from_lumina_maps_every_variant() {
338+
use crate::lumina::LuminaVectorMetric;
339+
assert_eq!(
340+
VectorSearchMetric::from_lumina(LuminaVectorMetric::L2),
341+
VectorSearchMetric::L2
342+
);
343+
assert_eq!(
344+
VectorSearchMetric::from_lumina(LuminaVectorMetric::Cosine),
345+
VectorSearchMetric::Cosine
346+
);
347+
assert_eq!(
348+
VectorSearchMetric::from_lumina(LuminaVectorMetric::InnerProduct),
349+
VectorSearchMetric::InnerProduct
350+
);
351+
}
323352
}

0 commit comments

Comments
 (0)