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 4131c4df perf(table): stream primary-key vector exact fallback one
batch at a time (#559)
4131c4df is described below
commit 4131c4dfc90e1e92460aa5c0be89e698d57b2d07
Author: Junrui Lee <[email protected]>
AuthorDate: Tue Jul 21 12:17:22 2026 +0800
perf(table): stream primary-key vector exact fallback one batch at a time
(#559)
---
.../paimon/src/table/pk_vector_data_file_reader.rs | 404 ++++++++++++-------
crates/paimon/src/table/pk_vector_orchestrator.rs | 424 ++++++++++++-------
crates/paimon/src/table/vector_search_builder.rs | 34 +-
crates/paimon/src/vindex/pkvector/bucket.rs | 447 ++++++++++++++-------
crates/paimon/src/vindex/pkvector/exact.rs | 100 +++--
crates/paimon/src/vindex/pkvector/reader.rs | 4 +
6 files changed, 929 insertions(+), 484 deletions(-)
diff --git a/crates/paimon/src/table/pk_vector_data_file_reader.rs
b/crates/paimon/src/table/pk_vector_data_file_reader.rs
index 33067b6a..d445f153 100644
--- a/crates/paimon/src/table/pk_vector_data_file_reader.rs
+++ b/crates/paimon/src/table/pk_vector_data_file_reader.rs
@@ -15,15 +15,20 @@
// specific language governing permissions and limitations
// under the License.
-//! Exact sequential vector reader over one data file's vector column. Mirrors
-//! Java `org.apache.paimon.index.pkvector.PkVectorDataFileReader`.
+//! Streaming exact vector search over one data file's vector column. Mirrors
+//! Java `org.apache.paimon.index.pkvector.PkVectorDataFileReader` +
+//! `PkVectorExactSearcher`.
//!
-//! The factory projects the single vector column, reads the whole file in
-//! physical order, and preloads it into memory as `Vec<Option<Vec<f32>>>`
-//! (a NULL row is `None`). Deletion vectors are deliberately NOT applied here:
-//! physical position must stay in lockstep with the segment ordinal so the
-//! bucket search can address rows by position. The returned reader then serves
-//! vectors from memory one physical row at a time.
+//! The factory projects the single vector column and, per uncovered file,
+//! streams the column one Arrow batch at a time — feeding each row into
per-query
+//! bounded Top-K heaps and dropping the batch — so peak memory is one batch
plus
+//! the heaps rather than the whole column. Deletion vectors and residual
filters
+//! are deliberately NOT applied by the stream itself: physical position must
stay
+//! in lockstep with the segment ordinal, so exclusion is folded in via the
+//! caller-supplied `is_excluded(position)` predicate. A NULL row is not
scored but
+//! still advances the physical position.
+
+use std::collections::BinaryHeap;
use arrow_array::{Array, FixedSizeListArray, Float32Array, ListArray};
use futures::TryStreamExt;
@@ -32,7 +37,9 @@ use crate::spec::{DataField, DataType};
use crate::table::data_file_reader::DataFileReader;
use crate::table::source::DataSplit;
use crate::vindex::pkvector::bucket::BucketActiveFile;
-use crate::vindex::pkvector::reader::PkVectorReader;
+use crate::vindex::pkvector::exact::{drain_best_first, push_bounded,
validate_query, WorstFirst};
+use crate::vindex::pkvector::metric::VectorSearchMetric;
+use crate::vindex::pkvector::result::PkVectorSearchResult;
fn data_invalid(message: impl Into<String>) -> crate::Error {
crate::Error::DataInvalid {
@@ -41,7 +48,8 @@ fn data_invalid(message: impl Into<String>) -> crate::Error {
}
}
-/// Builds an exact [`PkVectorReader`] over one data file's vector column.
+/// Runs a streaming exact [`PkVectorSearchResult`] search over one data file's
+/// vector column.
///
/// `reader` is configured (via [`DataFileReader::with_read_type`]) to project
/// only the vector column, so each read returns a single-column batch. Mirrors
@@ -79,14 +87,33 @@ impl DataFilePkVectorReaderFactory {
})
}
- /// Preload the whole vector column of `file` into memory and return a
- /// sequential reader over it. `file` must name a data file present in this
- /// factory's split. The drained row count is checked against the file's
- /// `DataFileMeta.row_count`.
- pub(crate) async fn create(
+ /// Stream the vector column of `file` one Arrow batch at a time and
return one
+ /// bounded, BEST_FIRST Top-K list per query (outer index aligned to
`queries`).
+ /// `file` must name a data file present in this factory's split.
+ ///
+ /// All queries are validated (dimension + finite) BEFORE the file stream
is
+ /// opened. Each surviving physical position (not NULL, not `is_excluded`)
is
+ /// scored against every query into that query's bounded heap; a NULL row
is
+ /// skipped but still advances the position so the position stays in
lockstep
+ /// with `is_excluded`. The drained row count is checked against the file's
+ /// `DataFileMeta.row_count` (both truncation and overrun fail loud).
+ pub(crate) async fn search_file(
&self,
file: &BucketActiveFile,
- ) -> crate::Result<Box<dyn PkVectorReader>> {
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync),
+ ) -> crate::Result<Vec<Vec<PkVectorSearchResult>>> {
+ if exact_limit == 0 {
+ return Err(data_invalid("vector search limit must be positive"));
+ }
+ // Validate every query before opening the stream
(validate-before-POLL);
+ // a malformed query fails loud before any file I/O.
+ for query in queries {
+ validate_query(query, self.dimension)?;
+ }
+
let file_meta = self
.data_split
.data_files()
@@ -110,34 +137,58 @@ impl DataFilePkVectorReaderFactory {
None,
)?;
- let mut vectors: Vec<Option<Vec<f32>>> = Vec::new();
+ let mut heaps: Vec<BinaryHeap<WorstFirst>> = (0..queries.len())
+ .map(|_| BinaryHeap::with_capacity(exact_limit + 1))
+ .collect();
+ // One reused buffer per batch; a NULL row leaves it untouched (and is
not
+ // scored). `position` is the monotonic physical row counter across
batches.
+ let mut batch_vectors: Vec<Option<Vec<f32>>> = Vec::new();
+ let mut position: i64 = 0;
while let Some(batch) = stream.try_next().await? {
+ batch_vectors.clear();
append_batch_vectors(
&batch,
self.vector_field.name(),
self.dimension,
- &mut vectors,
+ &mut batch_vectors,
)?;
+ for entry in &batch_vectors {
+ let pos = position;
+ position += 1;
+ if pos >= row_count {
+ return Err(data_invalid(
+ "data file produced more rows than
DataFileMeta.row_count",
+ ));
+ }
+ let Some(vector) = entry else {
+ continue; // NULL row: not scored, position already
advanced.
+ };
+ if is_excluded(pos) {
+ continue;
+ }
+ for (query, heap) in queries.iter().zip(heaps.iter_mut()) {
+ let candidate = PkVectorSearchResult {
+ data_file_name: file.file_name.clone(),
+ row_position: pos,
+ distance: metric.compute_distance(query, vector),
+ };
+ push_bounded(heap, candidate, exact_limit);
+ }
+ }
}
- let drained = vectors.len() as i64;
- if drained > row_count {
+ if position > row_count {
return Err(data_invalid(
"data file produced more rows than DataFileMeta.row_count",
));
}
- if drained < row_count {
+ if position < row_count {
return Err(data_invalid(
"data file ended before DataFileMeta.row_count",
));
}
- Ok(Box::new(DataFilePkVectorReader {
- dimension: self.dimension,
- row_count,
- vectors,
- position: 0,
- }))
+ Ok(heaps.into_iter().map(drain_best_first).collect())
}
}
@@ -219,97 +270,6 @@ pub(crate) fn append_batch_vectors(
Ok(())
}
-/// In-memory sequential reader over one file's preloaded vector column. Each
-/// [`read_next_vector`](PkVectorReader::read_next_vector) advances exactly one
-/// physical row; a NULL row returns `false` but still advances the position.
-struct DataFilePkVectorReader {
- dimension: usize,
- row_count: i64,
- /// Preloaded whole-file column in physical order; `None` = NULL row.
- vectors: Vec<Option<Vec<f32>>>,
- position: usize,
-}
-
-impl PkVectorReader for DataFilePkVectorReader {
- fn dimension(&self) -> usize {
- self.dimension
- }
-
- fn row_count(&self) -> i64 {
- self.row_count
- }
-
- fn read_next_vector(&mut self, reuse: &mut [f32]) -> crate::Result<bool> {
- if reuse.len() != self.dimension {
- return Err(data_invalid(format!(
- "reuse buffer length {} does not match vector dimension {}",
- reuse.len(),
- self.dimension
- )));
- }
- if self.position as i64 >= self.row_count {
- return Err(data_invalid("read past row count"));
- }
- let entry = &self.vectors[self.position];
- self.position += 1;
- match entry {
- Some(vector) => {
- reuse.copy_from_slice(vector);
- Ok(true)
- }
- None => Ok(false),
- }
- }
-}
-
-#[cfg(test)]
-mod tests {
- use super::*;
-
- #[test]
- fn null_row_consumes_ordinal() {
- let mut r = DataFilePkVectorReader {
- dimension: 2,
- row_count: 3,
- vectors: vec![Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.0])],
- position: 0,
- };
- let mut buf = [0.0f32; 2];
- assert!(r.read_next_vector(&mut buf).unwrap());
- assert_eq!(buf, [1.0, 2.0]);
- assert!(!r.read_next_vector(&mut buf).unwrap()); // null: false,
ordinal advanced
- assert!(r.read_next_vector(&mut buf).unwrap());
- assert_eq!(buf, [3.0, 4.0]);
- assert_eq!(r.row_count(), 3);
- assert_eq!(r.dimension(), 2);
- }
-
- #[test]
- fn read_past_row_count_errors() {
- let mut r = DataFilePkVectorReader {
- dimension: 1,
- row_count: 1,
- vectors: vec![Some(vec![1.0])],
- position: 0,
- };
- let mut buf = [0.0f32; 1];
- assert!(r.read_next_vector(&mut buf).unwrap());
- assert!(r.read_next_vector(&mut buf).is_err()); // past row_count
- }
-
- #[test]
- fn reuse_len_mismatch_errors() {
- let mut r = DataFilePkVectorReader {
- dimension: 2,
- row_count: 1,
- vectors: vec![Some(vec![1.0, 2.0])],
- position: 0,
- };
- let mut buf = [0.0f32; 1];
- assert!(r.read_next_vector(&mut buf).is_err());
- }
-}
-
#[cfg(test)]
mod integration_tests {
use super::*;
@@ -384,12 +344,16 @@ mod integration_tests {
);
}
- /// Write a FixedSizeList<Float32, 2> vector column
- /// (`[1,2]`, NULL, `[3,4]`) as a parquet data file, build the factory over
- /// its split, preload via `create`, and assert the whole-file sequential
- /// read plus the "file not in split" error.
- #[tokio::test]
- async fn create_preloads_and_reads_whole_file() {
+ /// Build a FixedSizeList<Float32, 2> vector column from `rows` (`None` =
NULL
+ /// row), write it as one parquet data file across `batches` write calls,
and
+ /// return a factory over its split plus the file name. `stated_row_count`
is
+ /// what the `DataFileMeta` claims (usually the true row count, but a test
can
+ /// pass a wrong value to exercise the row-count guard).
+ async fn build_factory(
+ rows: &[Option<Vec<f32>>],
+ stated_row_count: i64,
+ table_path: &str,
+ ) -> (DataFilePkVectorReaderFactory, String) {
let field = vector_field();
let read_fields = vec![field.clone()];
let arrow_schema = build_target_arrow_schema(&read_fields).unwrap();
@@ -397,20 +361,24 @@ mod integration_tests {
let mut builder = FixedSizeListBuilder::new(Float32Builder::new(),
2).with_field(Arc::new(
ArrowField::new("element", ArrowDataType::Float32, true),
));
- builder.values().append_value(1.0);
- builder.values().append_value(2.0);
- builder.append(true);
- builder.values().append_value(0.0);
- builder.values().append_value(0.0);
- builder.append(false); // NULL vector row
- builder.values().append_value(3.0);
- builder.values().append_value(4.0);
- builder.append(true);
+ for row in rows {
+ match row {
+ Some(v) => {
+ builder.values().append_value(v[0]);
+ builder.values().append_value(v[1]);
+ builder.append(true);
+ }
+ None => {
+ builder.values().append_value(0.0);
+ builder.values().append_value(0.0);
+ builder.append(false);
+ }
+ }
+ }
let vec_array = builder.finish();
let batch = RecordBatch::try_new(arrow_schema.clone(),
vec![Arc::new(vec_array)]).unwrap();
let file_io = FileIOBuilder::new("memory").build().unwrap();
- let table_path = "memory:/pk_vector_data_file_reader";
let bucket_path = format!("{table_path}/bucket-0");
let file_name = "part-0.parquet";
let file_path = format!("{bucket_path}/{file_name}");
@@ -440,7 +408,7 @@ mod integration_tests {
.with_data_files(vec![data_file(
file_name,
file_size as i64,
- 3,
+ stated_row_count,
table_schema_id,
)])
.build()
@@ -455,36 +423,162 @@ mod integration_tests {
read_fields.clone(),
Vec::new(),
);
+ let factory = DataFilePkVectorReaderFactory::new(reader, data_split,
field).unwrap();
+ (factory, file_name.to_string())
+ }
- let factory =
- DataFilePkVectorReaderFactory::new(reader, data_split,
field.clone()).unwrap();
+ /// The streaming per-file search must produce candidates byte-identical to
+ /// the reference `exact_search` over an in-memory `ArrayReader` of the
same
+ /// data, including a NULL row and a residual/DV exclusion.
+ #[tokio::test]
+ async fn search_file_matches_exact_search_reference() {
+ use crate::vindex::pkvector::exact::exact_search;
+ use crate::vindex::pkvector::reader::test_support::ArrayReader;
- let present = BucketActiveFile {
- file_name: file_name.to_string(),
+ let rows = vec![
+ Some(vec![3.0, 0.0]),
+ None,
+ Some(vec![1.0, 0.0]),
+ Some(vec![2.0, 0.0]),
+ ];
+ let (factory, file_name) =
+ build_factory(&rows, rows.len() as i64,
"memory:/pkvdfr_equiv").await;
+ let active = BucketActiveFile {
+ file_name: file_name.clone(),
+ row_count: rows.len() as i64,
+ };
+ // Exclude physical position 2 (residual/DV fold): mirrors the closure
the
+ // bucket search passes in.
+ let is_excluded = |pos: i64| pos == 2;
+ let query = [0.0f32, 0.0];
+
+ let streamed = factory
+ .search_file(&active, &[&query], VectorSearchMetric::L2, 2,
&is_excluded)
+ .await
+ .unwrap();
+
+ let mut ref_reader = ArrayReader::new(2, rows.clone());
+ let reference = exact_search(
+ &file_name,
+ &mut ref_reader,
+ &query,
+ VectorSearchMetric::L2,
+ 2,
+ &is_excluded,
+ )
+ .unwrap();
+
+ assert_eq!(streamed.len(), 1, "one query in, one result list out");
+ assert_eq!(streamed[0], reference);
+ }
+
+ /// Same streaming-vs-reference equivalence, but with more scorable rows
than
+ /// `exact_limit` so the bounded heap's eviction branch is exercised on
both
+ /// paths (the shared `push_bounded` must evict identically).
+ #[tokio::test]
+ async fn search_file_matches_exact_search_reference_with_eviction() {
+ use crate::vindex::pkvector::exact::exact_search;
+ use crate::vindex::pkvector::reader::test_support::ArrayReader;
+
+ // Five scorable rows, no NULL/exclusion; keep only the 2 closest to
[0,0].
+ let rows = vec![
+ Some(vec![4.0, 0.0]),
+ Some(vec![1.0, 0.0]),
+ Some(vec![3.0, 0.0]),
+ Some(vec![2.0, 0.0]),
+ Some(vec![5.0, 0.0]),
+ ];
+ let (factory, file_name) =
+ build_factory(&rows, rows.len() as i64,
"memory:/pkvdfr_evict").await;
+ let active = BucketActiveFile {
+ file_name: file_name.clone(),
+ row_count: rows.len() as i64,
+ };
+ let query = [0.0f32, 0.0];
+
+ let streamed = factory
+ .search_file(&active, &[&query], VectorSearchMetric::L2, 2, &|_|
false)
+ .await
+ .unwrap();
+
+ let mut ref_reader = ArrayReader::new(2, rows.clone());
+ let reference = exact_search(
+ &file_name,
+ &mut ref_reader,
+ &query,
+ VectorSearchMetric::L2,
+ 2,
+ &|_| false,
+ )
+ .unwrap();
+
+ assert_eq!(streamed[0], reference);
+ // The two closest are positions 1 ([1,0]) then 3 ([2,0]), best-first.
+ assert_eq!(streamed[0].len(), 2, "bounded to exact_limit");
+ assert_eq!(streamed[0][0].row_position, 1);
+ assert_eq!(streamed[0][1].row_position, 3);
+ }
+
+ /// A `DataFileMeta.row_count` larger than the file's real row count means
the
+ /// stream ends early; the search must fail loud rather than return a short
+ /// result.
+ #[tokio::test]
+ async fn search_file_fails_loud_on_row_count_truncation() {
+ let rows = vec![Some(vec![1.0, 0.0]), Some(vec![2.0, 0.0])];
+ // Claim 3 rows but only write 2.
+ let (factory, file_name) = build_factory(&rows, 3,
"memory:/pkvdfr_trunc").await;
+ let active = BucketActiveFile {
+ file_name,
row_count: 3,
};
- let mut pk_reader = factory.create(&present).await.unwrap();
- assert_eq!(pk_reader.dimension(), 2);
- assert_eq!(pk_reader.row_count(), 3);
-
- let mut buf = [0.0f32; 2];
- assert!(pk_reader.read_next_vector(&mut buf).unwrap());
- assert_eq!(buf, [1.0, 2.0]);
- assert!(!pk_reader.read_next_vector(&mut buf).unwrap()); // NULL row
- assert!(pk_reader.read_next_vector(&mut buf).unwrap());
- assert_eq!(buf, [3.0, 4.0]);
- assert!(pk_reader.read_next_vector(&mut buf).is_err()); // past row
count
-
- // A file name absent from the split is rejected as invalid.
+ let query = [0.0f32, 0.0];
+ let err = factory
+ .search_file(&active, &[&query], VectorSearchMetric::L2, 2, &|_|
false)
+ .await
+ .expect_err("row-count truncation must fail loud");
+ assert!(err.to_string().contains("ended before"), "got: {err}");
+ }
+
+ /// A malformed query (wrong dimension / non-finite element) fails loud,
and a
+ /// file name absent from the split is rejected as invalid.
+ #[tokio::test]
+ async fn search_file_validates_query_and_rejects_absent_file() {
+ let rows = vec![Some(vec![1.0, 2.0]), None, Some(vec![3.0, 4.0])];
+ let (factory, file_name) =
+ build_factory(&rows, rows.len() as i64,
"memory:/pkvdfr_validate").await;
+ let present = BucketActiveFile {
+ file_name: file_name.clone(),
+ row_count: rows.len() as i64,
+ };
+
+ // Wrong dimension.
+ let bad_dim = [1.0f32];
+ let err = factory
+ .search_file(&present, &[&bad_dim], VectorSearchMetric::L2, 2,
&|_| false)
+ .await
+ .expect_err("dimension mismatch must fail loud");
+ assert!(err.to_string().contains("dimension"), "got: {err}");
+
+ // Non-finite element.
+ let bad_finite = [f32::NAN, 0.0];
+ let err = factory
+ .search_file(&present, &[&bad_finite], VectorSearchMetric::L2, 2,
&|_| {
+ false
+ })
+ .await
+ .expect_err("non-finite query must fail loud");
+ assert!(err.to_string().contains("finite"), "got: {err}");
+
+ // Absent file.
let missing = BucketActiveFile {
file_name: "absent.parquet".to_string(),
row_count: 3,
};
+ let query = [0.0f32, 0.0];
let err = factory
- .create(&missing)
+ .search_file(&missing, &[&query], VectorSearchMetric::L2, 2, &|_|
false)
.await
- .err()
- .expect("absent file must be rejected");
+ .expect_err("absent file must be rejected");
assert!(matches!(err, crate::Error::DataInvalid { .. }));
}
}
diff --git a/crates/paimon/src/table/pk_vector_orchestrator.rs
b/crates/paimon/src/table/pk_vector_orchestrator.rs
index 358211a5..3e6a09fa 100644
--- a/crates/paimon/src/table/pk_vector_orchestrator.rs
+++ b/crates/paimon/src/table/pk_vector_orchestrator.rs
@@ -35,7 +35,7 @@ use
crate::table::pk_vector_indexed_split_read::PkVectorIndexedSplit;
use crate::table::source::{DataSplit, DataSplitBuilder, RowRange};
use crate::vindex::pkvector::ann::PkVectorAnnSearcher;
use crate::vindex::pkvector::bucket::{
- bucket_search, BucketActiveFile, BucketAnnSegment, ExactReaderFuture,
+ bucket_search, BucketActiveFile, BucketAnnSegment, ExactFileSearchFuture,
};
use crate::vindex::pkvector::metric::{java_float_compare, VectorSearchMetric};
use crate::vindex::pkvector::result::PkVectorSearchResult;
@@ -47,29 +47,44 @@ fn data_invalid(message: impl Into<String>) -> crate::Error
{
}
}
-/// Coerce a closure into the higher-ranked per-file exact-reader factory
shape so
-/// its returned future borrows for exactly the file argument's lifetime.
Closure
-/// return types cannot express this borrow through inference alone, so the
bound
-/// is supplied here.
-fn as_bucket_factory<F>(f: F) -> F
+/// Coerce a closure into the higher-ranked split-scoped exact-file search
shape
+/// expected by [`PkVectorOrchestrator::search_candidates`], binding the
returned
+/// future's borrow to the arguments' lifetime. Callers building a search
closure
+/// use this so the higher-ranked bound is supplied where inference cannot. The
+/// closure is `Fn + Send + Sync` (it is called concurrently by the parallel
+/// search), not `FnMut`.
+#[allow(clippy::type_complexity)]
+pub(crate) fn as_split_exact_file_search<F>(f: F) -> F
where
- F: for<'f> FnMut(&'f BucketActiveFile) -> ExactReaderFuture<'f> + Send,
+ F: for<'s, 'a> Fn(
+ usize,
+ &'s PkVectorSearchSplit,
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync,
{
f
}
-/// Coerce a closure into the split-scoped exact-reader factory shape expected
by
-/// [`PkVectorOrchestrator::search_candidates`], binding the returned future's
-/// borrow to the file argument's lifetime. Callers building a factory closure
use
-/// this so the higher-ranked bound is supplied where inference cannot.
-pub(crate) fn as_split_exact_reader_factory<F>(f: F) -> F
+/// Coerce a closure into the per-file exact-file search shape `bucket_search`
+/// expects, supplying the higher-ranked bound closure inference cannot
express.
+#[allow(clippy::type_complexity)]
+fn as_bucket_exact_file_search<F>(f: F) -> F
where
- F: for<'s, 'f> FnMut(
+ F: for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
usize,
- &'s PkVectorSearchSplit,
- &'f BucketActiveFile,
- ) -> ExactReaderFuture<'f>
- + Send,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync,
{
f
}
@@ -356,12 +371,17 @@ impl PkVectorOrchestrator {
limit: usize,
indexed_limit: usize,
ann_searcher: Option<&dyn PkVectorAnnSearcher>,
- exact_reader_factory: &mut (dyn for<'s, 'f> FnMut(
+ exact_file_search: &(dyn for<'s, 'a> Fn(
usize,
&'s PkVectorSearchSplit,
- &'f BucketActiveFile,
- ) -> ExactReaderFuture<'f>
- + Send),
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync),
search_options: &HashMap<String, String>,
skip_exact_fallback: bool,
residual_by_split: Option<&[HashMap<String, RoaringTreemap>]>,
@@ -389,19 +409,35 @@ impl PkVectorOrchestrator {
let mut exact_candidates: Vec<PkVectorCandidate> = Vec::new();
for (split_index, split) in splits.iter().enumerate() {
let dvs = build_bucket_dv_map(&self.reader, split).await?;
- // Wrap the split-scoped factory into bucket_search's per-file
signature.
- // The coercion helper ties the produced future's borrow to the
file
- // argument, which closure inference cannot express on its own.
- let mut bucket_factory = as_bucket_factory(|file:
&BucketActiveFile| {
- exact_reader_factory(split_index, split, file)
- });
+ // Adapt the split-scoped search closure to bucket_search's
per-file
+ // closure by binding the current split index/split. The coercion
helper
+ // ties the produced future's borrow to the arguments, which
closure
+ // inference cannot express on its own.
+ let bucket_search_closure = as_bucket_exact_file_search(
+ |file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ exact_file_search(
+ split_index,
+ split,
+ file,
+ queries,
+ metric,
+ exact_limit,
+ is_excluded,
+ )
+ },
+ );
let residual_ranges = residual_by_split.map(|per_split|
&per_split[split_index]);
let result = bucket_search(
ann_searcher,
&split.ann_segments,
&split.active_files,
&dvs,
- &mut bucket_factory,
+ &bucket_search_closure,
query,
metric,
indexed_limit,
@@ -729,8 +765,8 @@ mod e2e_tests {
use crate::table::pk_vector_position_read::{PKEY_VECTOR_POSITION_COLUMN,
SEARCH_SCORE_COLUMN};
use crate::table::schema_manager::SchemaManager;
use crate::table::source::DeletionFile;
+ use crate::vindex::pkvector::exact::exact_search;
use crate::vindex::pkvector::reader::test_support::ArrayReader;
- use crate::vindex::pkvector::reader::PkVectorReader;
use arrow_array::{Array, Float32Array, Int32Array, Int64Array,
RecordBatch};
use bytes::Bytes;
use futures::TryStreamExt;
@@ -897,30 +933,114 @@ mod e2e_tests {
}
}
- /// Coerce a closure into the higher-ranked split-scoped exact-reader
factory
- /// shape so its returned future borrows for exactly the file argument's
- /// lifetime. Closure return types cannot express this borrow through
inference
- /// alone, so the bound is supplied here.
- fn as_split_factory<F>(f: F) -> F
+ /// Coerce a closure into the split-scoped exact-file search shape,
supplying
+ /// the higher-ranked bound closure inference cannot express. The closure
is
+ /// `Fn + Send + Sync`.
+ #[allow(clippy::type_complexity)]
+ fn as_split_search<F>(f: F) -> F
where
- F: for<'s, 'f> FnMut(
+ F: for<'s, 'a> Fn(
usize,
&'s PkVectorSearchSplit,
- &'f BucketActiveFile,
- ) -> ExactReaderFuture<'f>
- + Send,
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync,
{
- as_split_exact_reader_factory(f)
+ as_split_exact_file_search(f)
+ }
+
+ /// A split-scoped search closure that must never be invoked.
+ #[allow(clippy::type_complexity)]
+ fn unreachable_split_search() -> impl for<'s, 'a> Fn(
+ usize,
+ &'s PkVectorSearchSplit,
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync {
+ as_split_search(
+ |_: usize,
+ _: &PkVectorSearchSplit,
+ _: &BucketActiveFile,
+ _: &[&[f32]],
+ _: VectorSearchMetric,
+ _: usize,
+ _: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ Box::pin(async { unreachable!("closure must not be invoked in
this test") })
+ },
+ )
}
- /// Same coercion for the per-file factory used by
`materialize_via_splits`.
- fn as_file_factory<F>(f: F) -> F
+ #[allow(clippy::type_complexity)]
+ fn as_file_search<F>(f: F) -> F
where
- F: for<'f> FnMut(&'f BucketActiveFile) -> ExactReaderFuture<'f> + Send,
+ F: for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync,
{
f
}
+ /// Build a per-file search closure that runs the reference `exact_search`
over
+ /// an in-memory `ArrayReader` whose vectors come from
`vectors_for(file_name)`.
+ /// The caller wires a single query, so the returned outer `Vec` has one
element.
+ #[allow(clippy::type_complexity)]
+ fn file_array_search(
+ vectors_for: fn(&str) -> Vec<Option<Vec<f32>>>,
+ ) -> impl for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync {
+ as_file_search(
+ move |file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ let vectors = vectors_for(&file.file_name);
+ let dimension = vectors
+ .first()
+ .and_then(|v| v.as_ref())
+ .map_or(2, |v| v.len());
+ let file_name = file.file_name.clone();
+ let query = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader = ArrayReader::new(dimension, vectors);
+ Ok(vec![exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?])
+ })
+ },
+ )
+ }
+
fn column_by_name<'a>(batch: &'a RecordBatch, name: &str) -> Option<&'a
Arc<dyn Array>> {
batch
.schema()
@@ -1000,6 +1120,7 @@ mod e2e_tests {
/// materialization path the production best-first read reorders on top
of; the
/// tests below drive it directly through its `pub(crate)` components.
#[allow(clippy::too_many_arguments)]
+ #[allow(clippy::type_complexity)]
async fn materialize_via_splits(
reader: DataFileReader,
splits: &[PkVectorSearchSplit],
@@ -1007,26 +1128,35 @@ mod e2e_tests {
metric: VectorSearchMetric,
limit: usize,
ann: Option<&dyn PkVectorAnnSearcher>,
- factory: &mut (dyn for<'f> FnMut(&'f BucketActiveFile) ->
ExactReaderFuture<'f> + Send),
+ search: &(dyn for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync),
opts: &HashMap<String, String>,
) -> crate::Result<Vec<RecordBatch>> {
let orch = PkVectorOrchestrator::new(reader.clone());
- // Wrap the per-file factory into the split-scoped shape
search_candidates
+ // Wrap the per-file search into the split-scoped shape
search_candidates
// expects; the split index/split are unused here.
- let mut wrapped =
- as_split_factory(|_: usize, _: &PkVectorSearchSplit, f:
&BucketActiveFile| factory(f));
+ let wrapped = as_split_search(
+ |_: usize,
+ _: &PkVectorSearchSplit,
+ file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ search(file, queries, metric, exact_limit, is_excluded)
+ },
+ );
let result = orch
.search_candidates(
- splits,
- query,
- metric,
- limit,
- limit,
- ann,
- &mut wrapped,
- opts,
- false,
- None,
+ splits, query, metric, limit, limit, ann, &wrapped, opts,
false, None,
)
.await?;
// Merge the two bounded lists into the best-first survivors the
@@ -1049,11 +1179,7 @@ mod e2e_tests {
let file_io = FileIOBuilder::new("memory").build().unwrap();
let reader = make_reader(file_io, "memory:/pkvo_zero");
let splits: Vec<PkVectorSearchSplit> = Vec::new();
- let mut factory = as_split_factory(
- |_: usize, _: &PkVectorSearchSplit, _: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!("no bucket search on
eager-rejected input") })
- },
- );
+ let factory = unreachable_split_search();
let opts = HashMap::new();
let err = PkVectorOrchestrator::new(reader)
.search_candidates(
@@ -1063,7 +1189,7 @@ mod e2e_tests {
0,
0,
None,
- &mut factory,
+ &factory,
&opts,
false,
None,
@@ -1079,11 +1205,7 @@ mod e2e_tests {
let file_io = FileIOBuilder::new("memory").build().unwrap();
let reader = make_reader(file_io, "memory:/pkvo_empty_query");
let splits: Vec<PkVectorSearchSplit> = Vec::new();
- let mut factory = as_split_factory(
- |_: usize, _: &PkVectorSearchSplit, _: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!("no bucket search on
eager-rejected input") })
- },
- );
+ let factory = unreachable_split_search();
let opts = HashMap::new();
let err = PkVectorOrchestrator::new(reader)
.search_candidates(
@@ -1093,7 +1215,7 @@ mod e2e_tests {
5,
5,
None,
- &mut factory,
+ &factory,
&opts,
false,
None,
@@ -1138,14 +1260,9 @@ mod e2e_tests {
}],
};
// Exact fallback scans only "exact.mosaic": pos0 {1,0} d=1.0, pos1
{2,0} d=4.0.
- let mut factory = as_file_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- let vectors = match f.file_name.as_str() {
- "exact.mosaic" => vec![Some(vec![1.0, 0.0]), Some(vec![2.0,
0.0])],
- other => panic!("unexpected exact scan of covered file
{other}"),
- };
- Box::pin(async move {
- Ok(Box::new(ArrayReader::new(2, vectors)) as Box<dyn
PkVectorReader>)
- })
+ let factory = file_array_search(|file_name| match file_name {
+ "exact.mosaic" => vec![Some(vec![1.0, 0.0]), Some(vec![2.0, 0.0])],
+ other => panic!("unexpected exact scan of covered file {other}"),
});
let opts = HashMap::new();
let batches = materialize_via_splits(
@@ -1155,7 +1272,7 @@ mod e2e_tests {
VectorSearchMetric::L2,
3,
Some(&ann),
- &mut factory,
+ &factory,
&opts,
)
.await
@@ -1221,23 +1338,18 @@ mod e2e_tests {
// b0: x = 1,4,6 -> d = 1,16,36. b1: x = 2,3,5 -> d = 4,9,25.
// Global best 3: d1 (b0 pos0 id10), d4 (b1 pos0 id20), d9 (b1 pos1
id21).
- let mut factory = as_file_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- let vectors = match f.file_name.as_str() {
- "b0.mosaic" => vec![
- Some(vec![1.0, 0.0]),
- Some(vec![4.0, 0.0]),
- Some(vec![6.0, 0.0]),
- ],
- "b1.mosaic" => vec![
- Some(vec![2.0, 0.0]),
- Some(vec![3.0, 0.0]),
- Some(vec![5.0, 0.0]),
- ],
- other => panic!("unexpected file {other}"),
- };
- Box::pin(async move {
- Ok(Box::new(ArrayReader::new(2, vectors)) as Box<dyn
PkVectorReader>)
- })
+ let factory = file_array_search(|file_name| match file_name {
+ "b0.mosaic" => vec![
+ Some(vec![1.0, 0.0]),
+ Some(vec![4.0, 0.0]),
+ Some(vec![6.0, 0.0]),
+ ],
+ "b1.mosaic" => vec![
+ Some(vec![2.0, 0.0]),
+ Some(vec![3.0, 0.0]),
+ Some(vec![5.0, 0.0]),
+ ],
+ other => panic!("unexpected file {other}"),
});
let opts = HashMap::new();
let batches = materialize_via_splits(
@@ -1247,7 +1359,7 @@ mod e2e_tests {
VectorSearchMetric::L2,
3,
None,
- &mut factory,
+ &factory,
&opts,
)
.await
@@ -1292,19 +1404,14 @@ mod e2e_tests {
};
// pos0 {1,0} d=1, pos1 {2,0} d=4 (DELETED), pos2 {3,0} d=9, pos3
{0,0} d=0.
- let mut factory = as_file_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- let vectors = match f.file_name.as_str() {
- "d.mosaic" => vec![
- Some(vec![1.0, 0.0]),
- Some(vec![2.0, 0.0]),
- Some(vec![3.0, 0.0]),
- Some(vec![0.0, 0.0]),
- ],
- other => panic!("unexpected file {other}"),
- };
- Box::pin(async move {
- Ok(Box::new(ArrayReader::new(2, vectors)) as Box<dyn
PkVectorReader>)
- })
+ let factory = file_array_search(|file_name| match file_name {
+ "d.mosaic" => vec![
+ Some(vec![1.0, 0.0]),
+ Some(vec![2.0, 0.0]),
+ Some(vec![3.0, 0.0]),
+ Some(vec![0.0, 0.0]),
+ ],
+ other => panic!("unexpected file {other}"),
});
let opts = HashMap::new();
let batches = materialize_via_splits(
@@ -1314,7 +1421,7 @@ mod e2e_tests {
VectorSearchMetric::L2,
4,
None,
- &mut factory,
+ &factory,
&opts,
)
.await
@@ -1355,18 +1462,13 @@ mod e2e_tests {
};
// pos0 {3,0} d=9, pos1 {1,0} d=1, pos2 {2,0} d=4. Best-first =
[1,2,0].
- let mut factory = as_file_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- let vectors = match f.file_name.as_str() {
- "o.mosaic" => vec![
- Some(vec![3.0, 0.0]),
- Some(vec![1.0, 0.0]),
- Some(vec![2.0, 0.0]),
- ],
- other => panic!("unexpected file {other}"),
- };
- Box::pin(async move {
- Ok(Box::new(ArrayReader::new(2, vectors)) as Box<dyn
PkVectorReader>)
- })
+ let factory = file_array_search(|file_name| match file_name {
+ "o.mosaic" => vec![
+ Some(vec![3.0, 0.0]),
+ Some(vec![1.0, 0.0]),
+ Some(vec![2.0, 0.0]),
+ ],
+ other => panic!("unexpected file {other}"),
});
let opts = HashMap::new();
let batches = materialize_via_splits(
@@ -1376,7 +1478,7 @@ mod e2e_tests {
VectorSearchMetric::L2,
3,
None,
- &mut factory,
+ &factory,
&opts,
)
.await
@@ -1416,18 +1518,35 @@ mod e2e_tests {
active_files: vec![active("c.mosaic", 3)],
};
// pos0 {3,0} d=9, pos1 {1,0} d=1, pos2 {2,0} d=4.
- let mut factory = as_split_factory(
- |_: usize, _: &PkVectorSearchSplit, f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- assert_eq!(f.file_name, "c.mosaic");
- Box::pin(async {
- Ok(Box::new(ArrayReader::new(
+ let factory = as_split_search(
+ |_: usize,
+ _: &PkVectorSearchSplit,
+ file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ assert_eq!(file.file_name, "c.mosaic");
+ let file_name = file.file_name.clone();
+ let query = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader = ArrayReader::new(
2,
vec![
Some(vec![3.0, 0.0]),
Some(vec![1.0, 0.0]),
Some(vec![2.0, 0.0]),
],
- )) as Box<dyn PkVectorReader>)
+ );
+ Ok(vec![exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?])
})
},
);
@@ -1440,7 +1559,7 @@ mod e2e_tests {
2,
2,
None,
- &mut factory,
+ &factory,
&opts,
false,
None,
@@ -1485,18 +1604,35 @@ mod e2e_tests {
active_files: vec![active("r.mosaic", 3)],
};
// pos0 {3,0} d=9, pos1 {1,0} d=1, pos2 {2,0} d=4.
- let mut factory = as_split_factory(
- |_: usize, _: &PkVectorSearchSplit, f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- assert_eq!(f.file_name, "r.mosaic");
- Box::pin(async {
- Ok(Box::new(ArrayReader::new(
+ let factory = as_split_search(
+ |_: usize,
+ _: &PkVectorSearchSplit,
+ file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ assert_eq!(file.file_name, "r.mosaic");
+ let file_name = file.file_name.clone();
+ let query = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader = ArrayReader::new(
2,
vec![
Some(vec![3.0, 0.0]),
Some(vec![1.0, 0.0]),
Some(vec![2.0, 0.0]),
],
- )) as Box<dyn PkVectorReader>)
+ );
+ Ok(vec![exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?])
})
},
);
@@ -1515,7 +1651,7 @@ mod e2e_tests {
3,
3,
None,
- &mut factory,
+ &factory,
&opts,
false,
Some(&residual_by_split),
@@ -1556,11 +1692,7 @@ mod e2e_tests {
ann_segments: Vec::new(),
active_files: vec![active("m.mosaic", 2)],
};
- let mut factory = as_split_factory(
- |_: usize, _: &PkVectorSearchSplit, _: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!("length guard must fire before
any bucket search") })
- },
- );
+ let factory = unreachable_split_search();
// Two residual maps for a single split.
let residual_by_split: Vec<HashMap<String, RoaringTreemap>> =
vec![HashMap::new(), HashMap::new()];
@@ -1573,7 +1705,7 @@ mod e2e_tests {
3,
3,
None,
- &mut factory,
+ &factory,
&opts,
false,
Some(&residual_by_split),
@@ -1606,11 +1738,7 @@ mod e2e_tests {
ann_segments: Vec::new(),
active_files: vec![active("f.mosaic", 2)],
};
- let mut factory = as_split_factory(
- |_: usize, _: &PkVectorSearchSplit, _: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!("fast mode must not read exact")
})
- },
- );
+ let factory = unreachable_split_search();
let opts = HashMap::new();
let result = PkVectorOrchestrator::new(make_reader(file_io,
table_path))
.search_candidates(
@@ -1620,7 +1748,7 @@ mod e2e_tests {
2,
2,
None,
- &mut factory,
+ &factory,
&opts,
true,
None,
diff --git a/crates/paimon/src/table/vector_search_builder.rs
b/crates/paimon/src/table/vector_search_builder.rs
index 9f50ee46..6f4cb27a 100644
--- a/crates/paimon/src/table/vector_search_builder.rs
+++ b/crates/paimon/src/table/vector_search_builder.rs
@@ -34,8 +34,8 @@ use crate::table::pk_vector_data_file_reader::{
};
use crate::table::pk_vector_indexed_split_read::{expand_ranges,
PkVectorIndexedSplitRead};
use crate::table::pk_vector_orchestrator::{
- as_split_exact_reader_factory, build_indexed_splits, merge_candidates,
- OrchestratorSearchResult, PkVectorCandidate, PkVectorOrchestrator,
PkVectorSearchSplit,
+ as_split_exact_file_search, build_indexed_splits, merge_candidates,
OrchestratorSearchResult,
+ PkVectorCandidate, PkVectorOrchestrator, PkVectorSearchSplit,
};
use crate::table::pk_vector_position_read::{
PkVectorPositionRead, PKEY_VECTOR_POSITION_COLUMN, SEARCH_SCORE_COLUMN,
@@ -49,7 +49,7 @@ use crate::table::{
use crate::vector_search::{GlobalIndexIOMeta, SearchResult, VectorSearch};
use crate::vindex::is_vindex_index_type;
use crate::vindex::pkvector::ann::VindexAnnSearcher;
-use crate::vindex::pkvector::bucket::{BucketActiveFile, BucketAnnSegment,
ExactReaderFuture};
+use crate::vindex::pkvector::bucket::{BucketActiveFile, BucketAnnSegment,
ExactFileSearchFuture};
use crate::vindex::pkvector::metric::VectorSearchMetric;
use crate::vindex::reader::VindexVectorGlobalIndexReader;
use arrow_array::{Array, FixedSizeListArray, Float32Array, Int64Array,
ListArray, RecordBatch};
@@ -577,18 +577,23 @@ impl<'a> VectorSearchBuilder<'a> {
None => None,
};
- // Build the exact-fallback vector reader on demand: the kernel calls
this
- // only for a file it actually searches (uncovered by ANN,
residual-allowed,
- // and only when the search mode is not FAST). Everything the future
needs is
+ // Build the exact-fallback search on demand: the kernel calls this
only
+ // for a file it actually searches (uncovered by ANN,
residual-allowed, and
+ // only when the search mode is not FAST). Everything the future needs
is
// cloned/owned up front so it borrows neither the split nor the file
across
- // the await.
+ // the await. The search streams the file's vector column one Arrow
batch at
+ // a time into per-query bounded heaps (the caller passes a single
query).
let reader_for_factory = reader.clone();
let vector_field_for_factory = vector_field.clone();
- let mut factory = as_split_exact_reader_factory(
+ let factory = as_split_exact_file_search(
move |_split_index: usize,
split: &PkVectorSearchSplit,
- file: &BucketActiveFile|
- -> ExactReaderFuture<'_> {
+ file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
let reader = reader_for_factory.clone();
let vector_field = vector_field_for_factory.clone();
let data_split = split.data_split.clone();
@@ -596,10 +601,15 @@ impl<'a> VectorSearchBuilder<'a> {
file_name: file.file_name.clone(),
row_count: file.row_count,
};
+ let owned_queries: Vec<Vec<f32>> = queries.iter().map(|q|
q.to_vec()).collect();
Box::pin(async move {
let factory =
DataFilePkVectorReaderFactory::new(reader, data_split,
vector_field)?;
- factory.create(&active).await
+ let query_refs: Vec<&[f32]> =
+ owned_queries.iter().map(|q| q.as_slice()).collect();
+ factory
+ .search_file(&active, &query_refs, metric,
exact_limit, is_excluded)
+ .await
})
},
);
@@ -612,7 +622,7 @@ impl<'a> VectorSearchBuilder<'a> {
limit,
indexed_limit,
Some(&ann_searcher),
- &mut factory,
+ &factory,
&search_options,
skip_exact_fallback,
residual_by_split.as_deref(),
diff --git a/crates/paimon/src/vindex/pkvector/bucket.rs
b/crates/paimon/src/vindex/pkvector/bucket.rs
index f3cfc35f..3b2f7fb1 100644
--- a/crates/paimon/src/vindex/pkvector/bucket.rs
+++ b/crates/paimon/src/vindex/pkvector/bucket.rs
@@ -23,16 +23,19 @@ use futures::future::BoxFuture;
use super::ann::PkVectorAnnSearcher;
use super::data_invalid;
-use super::exact::exact_search;
use super::metric::{java_float_compare, VectorSearchMetric};
-use super::reader::PkVectorReader;
use super::result::PkVectorSearchResult;
use crate::deletion_vector::DeletionVector;
use crate::spec::PkVectorSourceMeta;
-/// Future returned by the exact-fallback reader factory: builds the sequential
-/// vector reader for one data file on demand.
-pub(crate) type ExactReaderFuture<'a> = BoxFuture<'a, crate::Result<Box<dyn
PkVectorReader>>>;
+/// Search one uncovered data file for its per-query exact Top-K. Returns one
+/// bounded, BEST_FIRST list per query (outer index aligns to the `queries`
slice
+/// passed to the closure). The table-layer implementation streams the file's
+/// vector column one Arrow batch at a time, so peak memory is one batch plus
the
+/// bounded heaps, not the whole column. `is_excluded(position)` folds
residual +
+/// deletion-vector exclusion (built by `bucket_search`, which owns those
inputs).
+pub(crate) type ExactFileSearchFuture<'a> =
+ BoxFuture<'a, crate::Result<Vec<Vec<PkVectorSearchResult>>>>;
/// One ANN segment to be searched by the bucket kernel. `source_meta` resolves
/// segment ordinals back to physical `(data file, position)` and drives
live-row
@@ -164,13 +167,21 @@ pub(crate) struct BucketSearchResult {
/// with an empty set) has no allowed rows and produces no candidates. Mirrors
Java
/// `rowRangesByFile`.
#[allow(clippy::too_many_arguments)]
+#[allow(clippy::type_complexity)]
pub(crate) async fn bucket_search(
ann_searcher: Option<&dyn PkVectorAnnSearcher>,
ann_segments: &[BucketAnnSegment],
active_files: &[BucketActiveFile],
deletion_vectors: &HashMap<String, Arc<DeletionVector>>,
- exact_reader_factory: &mut (dyn for<'a> FnMut(&'a BucketActiveFile) ->
ExactReaderFuture<'a>
- + Send),
+ exact_file_search: &(dyn for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync),
query: &[f32],
metric: VectorSearchMetric,
indexed_limit: usize,
@@ -185,6 +196,15 @@ pub(crate) async fn bucket_search(
if exact_limit == 0 {
return Err(data_invalid("vector search limit must be positive"));
}
+ // Validate-before-OPEN: a non-finite query element fails loud before any
+ // exact-file search closure is invoked (so before any file stream is
opened).
+ // The per-file closure additionally validates the query dimension before
it
+ // polls its stream (validate-before-POLL).
+ if let Some(i) = query.iter().position(|v| !v.is_finite()) {
+ return Err(data_invalid(format!(
+ "query vector element at position {i} must be finite"
+ )));
+ }
let mut files_by_name: HashMap<&str, &BucketActiveFile> = HashMap::new();
for file in active_files {
@@ -307,15 +327,22 @@ pub(crate) async fn bucket_search(
},
}
};
- let mut reader = exact_reader_factory(file).await?;
- for result in exact_search(
- &file.file_name,
- reader.as_mut(),
- query,
+ // Search this one file for its per-query exact Top-K. The caller
passes
+ // a single-query slice and reads the sole returned list.
+ let queries: [&[f32]; 1] = [query];
+ let per_query = exact_file_search(
+ file,
+ &queries,
metric,
exact_limit,
- &is_excluded,
- )? {
+ &is_excluded as &(dyn Fn(i64) -> bool + Sync),
+ )
+ .await?;
+ let results = per_query
+ .into_iter()
+ .next()
+ .ok_or_else(|| data_invalid("exact file search returned no
per-query results"))?;
+ for result in results {
add_candidate(&mut exact_heap, result, exact_limit);
}
}
@@ -333,6 +360,7 @@ mod tests {
use super::*;
use crate::spec::PkVectorSourceFile;
use crate::vindex::pkvector::ann::PkVectorAnnSearcher;
+ use crate::vindex::pkvector::exact::exact_search;
use crate::vindex::pkvector::reader::test_support::ArrayReader;
use roaring::RoaringBitmap;
@@ -354,17 +382,93 @@ mod tests {
}
}
- /// Coerce a closure into the higher-ranked exact-reader factory shape so
its
- /// returned future borrows for exactly the argument's lifetime. Closure
return
+ /// Coerce a closure into the higher-ranked per-file exact-search shape so
its
+ /// returned future borrows for exactly the argument lifetime. Closure
return
/// types cannot express this borrow through inference alone, so the bound
is
- /// supplied here.
- fn as_factory<F>(f: F) -> F
+ /// supplied here. The closure is `Fn + Send + Sync` (called concurrently
by
+ /// the parallel search), not `FnMut`.
+ fn as_search<F>(f: F) -> F
where
- F: for<'a> FnMut(&'a BucketActiveFile) -> ExactReaderFuture<'a> + Send,
+ F: for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync,
{
f
}
+ /// Build a per-file exact-search closure that runs the reference
`exact_search`
+ /// over an in-memory `ArrayReader` of `vectors`, returning one per-query
list.
+ /// The caller wires a single query, so the returned outer `Vec` has one
element.
+ #[allow(clippy::type_complexity)]
+ fn array_search(
+ vectors: Vec<Option<Vec<f32>>>,
+ ) -> impl for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync {
+ let vectors = Arc::new(vectors);
+ as_search(
+ move |file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ let dimension = vectors
+ .first()
+ .and_then(|v| v.as_ref())
+ .map_or(0, |v| v.len());
+ let owned: Vec<Option<Vec<f32>>> = (*vectors).clone();
+ let file_name = file.file_name.clone();
+ let query: Vec<f32> = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader = ArrayReader::new(dimension, owned);
+ let results = exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?;
+ Ok(vec![results])
+ })
+ },
+ )
+ }
+
+ /// A per-file exact-search closure that must never be invoked.
+ #[allow(clippy::type_complexity)]
+ fn unreachable_search() -> impl for<'a> Fn(
+ &'a BucketActiveFile,
+ &'a [&'a [f32]],
+ VectorSearchMetric,
+ usize,
+ &'a (dyn Fn(i64) -> bool + Sync),
+ ) -> ExactFileSearchFuture<'a>
+ + Send
+ + Sync {
+ as_search(
+ |_: &BucketActiveFile,
+ _: &[&[f32]],
+ _: VectorSearchMetric,
+ _: usize,
+ _: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> { Box::pin(async { unreachable!() })
},
+ )
+ }
+
/// Fake ANN searcher returning preset results and recording calls.
struct FakeAnnSearcher {
result: Vec<PkVectorSearchResult>,
@@ -412,17 +516,11 @@ mod tests {
],
};
// Exact file has three rows; query nearest is position 0.
- let mut factory = as_factory(|_f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- let reader = ArrayReader::new(
- 2,
- vec![
- Some(vec![1.0, 0.0]),
- Some(vec![9.0, 0.0]),
- Some(vec![8.0, 0.0]),
- ],
- );
- Box::pin(async move { Ok(Box::new(reader) as Box<dyn
PkVectorReader>) })
- });
+ let factory = array_search(vec![
+ Some(vec![1.0, 0.0]),
+ Some(vec![9.0, 0.0]),
+ Some(vec![8.0, 0.0]),
+ ]);
let active_files = vec![active("ann.mosaic", 3),
active("exact.mosaic", 3)];
let dvs: HashMap<String, Arc<DeletionVector>> = HashMap::new();
let opts = HashMap::new();
@@ -432,7 +530,7 @@ mod tests {
&[segment],
&active_files,
&dvs,
- &mut factory,
+ &factory,
&[1.0, 0.0],
VectorSearchMetric::L2,
3, // indexed_limit
@@ -456,15 +554,13 @@ mod tests {
#[tokio::test]
async fn test_rejects_non_positive_limit() {
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
None,
&[],
&[],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
0,
@@ -500,15 +596,13 @@ mod tests {
hit("data-1", 1),
],
};
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let out = bucket_search(
Some(&ann),
&[segment],
&[active("data-1", 3)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
3,
@@ -557,15 +651,13 @@ mod tests {
},
],
};
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let out = bucket_search(
Some(&ann),
&[segment],
&[active("data-1", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -598,22 +690,37 @@ mod tests {
}],
};
let calls = std::sync::Mutex::new(Vec::<String>::new());
- let mut factory = as_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- calls.lock().unwrap().push(f.file_name.clone());
- // data-2 vectors: pos0 {1,0} dist 1.0, pos1 {3,0} dist 9.0
- Box::pin(async {
- Ok(Box::new(ArrayReader::new(
- 2,
- vec![Some(vec![1.0, 0.0]), Some(vec![3.0, 0.0])],
- )) as Box<dyn PkVectorReader>)
- })
- });
+ let factory = as_search(
+ |file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ calls.lock().unwrap().push(file.file_name.clone());
+ // data-2 vectors: pos0 {1,0} dist 1.0, pos1 {3,0} dist 9.0
+ let file_name = file.file_name.clone();
+ let query = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader =
+ ArrayReader::new(2, vec![Some(vec![1.0, 0.0]),
Some(vec![3.0, 0.0])]);
+ Ok(vec![exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?])
+ })
+ },
+ );
let out = bucket_search(
Some(&ann),
&[segment],
&[active("data-1", 2), active("data-2", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
2,
@@ -650,17 +757,34 @@ mod tests {
async fn test_exact_fallback_merges_files_and_applies_deletion_vectors() {
// No ANN. data-1 pos0 {0,0} deleted; remaining candidates merge
across files.
let calls = std::sync::Mutex::new(0usize);
- let mut factory = as_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- *calls.lock().unwrap() += 1;
- let vectors = match f.file_name.as_str() {
- "data-1" => vec![Some(vec![0.0, 0.0]), Some(vec![2.0, 0.0])],
- "data-2" => vec![Some(vec![1.0, 0.0]), None],
- _ => unreachable!(),
- };
- Box::pin(async move {
- Ok(Box::new(ArrayReader::new(2, vectors)) as Box<dyn
PkVectorReader>)
- })
- });
+ let factory = as_search(
+ |file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ *calls.lock().unwrap() += 1;
+ let vectors = match file.file_name.as_str() {
+ "data-1" => vec![Some(vec![0.0, 0.0]), Some(vec![2.0,
0.0])],
+ "data-2" => vec![Some(vec![1.0, 0.0]), None],
+ _ => unreachable!(),
+ };
+ let file_name = file.file_name.clone();
+ let query = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader = ArrayReader::new(2, vectors);
+ Ok(vec![exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?])
+ })
+ },
+ );
let mut dvs: HashMap<String, Arc<DeletionVector>> = HashMap::new();
let mut bm = RoaringBitmap::new();
bm.insert(0); // data-1 position 0 deleted
@@ -671,7 +795,7 @@ mod tests {
&[],
&[active("data-1", 2), active("data-2", 2)],
&dvs,
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
2,
@@ -707,15 +831,13 @@ mod tests {
#[tokio::test]
async fn test_rejects_duplicate_active_file_name() {
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
None,
&[],
&[active("dup", 1), active("dup", 1)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -735,15 +857,13 @@ mod tests {
// Segment references data-1 with 2 rows, but the active file has 3
rows.
// An active source with a mismatched row count is still a hard error.
let segment = BucketAnnSegment::for_test(meta(&[("data-1", 2)]));
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
Some(&ann),
&[segment],
&[active("data-1", 3)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -775,16 +895,23 @@ mod tests {
}],
};
let calls = std::sync::Mutex::new(Vec::<String>::new());
- let mut factory = as_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- calls.lock().unwrap().push(f.file_name.clone());
- Box::pin(async { unreachable!("only data-1 is active and it is
ANN-covered") })
- });
+ let factory = as_search(
+ |file: &BucketActiveFile,
+ _: &[&[f32]],
+ _: VectorSearchMetric,
+ _: usize,
+ _: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ calls.lock().unwrap().push(file.file_name.clone());
+ Box::pin(async { unreachable!("only data-1 is active and it is
ANN-covered") })
+ },
+ );
let out = bucket_search(
Some(&ann),
&[segment],
&[active("data-1", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
2,
@@ -814,15 +941,13 @@ mod tests {
#[tokio::test]
async fn test_rejects_segments_without_ann_searcher() {
let segment = BucketAnnSegment::for_test(meta(&[("data-1", 2)]));
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
None,
&[segment],
&[active("data-1", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -843,15 +968,13 @@ mod tests {
async fn test_skip_exact_fallback_does_not_call_factory() {
// No ANN segments, two active files. With skip_exact_fallback = true
the
// factory must never be called and the result is empty.
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let out = bucket_search(
None,
&[],
&[active("data-1", 2), active("data-2", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
2,
@@ -881,15 +1004,13 @@ mod tests {
index_meta: vec![4, 5, 6],
};
let ann = FakeAnnSearcher { result: vec![] };
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
Some(&ann),
&[seg1, seg2],
&[active("data-1", 2), active("data-2", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -921,15 +1042,13 @@ mod tests {
index_meta: vec![4, 5, 6],
};
let ann = FakeAnnSearcher { result: vec![] };
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
Some(&ann),
&[seg1, seg2],
&[active("data-1", 2), active("data-2", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -950,15 +1069,13 @@ mod tests {
#[tokio::test]
async fn test_negative_active_row_count_rejected() {
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async { unreachable!() })
- });
+ let factory = unreachable_search();
let err = bucket_search(
None,
&[],
&[active("data-1", -1)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
1,
@@ -1014,18 +1131,11 @@ mod tests {
// No ANN. data-1 has 3 rows: pos0 {1,0} dist 1.0, pos1 {2,0} dist 4.0,
// pos2 {3,0} dist 9.0. residual allows only {0, 2} -> pos1 excluded
even
// though it is not deletion-vector deleted.
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async {
- Ok(Box::new(ArrayReader::new(
- 2,
- vec![
- Some(vec![1.0, 0.0]),
- Some(vec![2.0, 0.0]),
- Some(vec![3.0, 0.0]),
- ],
- )) as Box<dyn PkVectorReader>)
- })
- });
+ let factory = array_search(vec![
+ Some(vec![1.0, 0.0]),
+ Some(vec![2.0, 0.0]),
+ Some(vec![3.0, 0.0]),
+ ]);
let mut residual: HashMap<String, roaring::RoaringTreemap> =
HashMap::new();
residual.insert("data-1".into(), treemap(&[0, 2]));
let out = bucket_search(
@@ -1033,7 +1143,7 @@ mod tests {
&[],
&[active("data-1", 3)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
5,
@@ -1070,15 +1180,30 @@ mod tests {
// residual covers only data-1; data-2 has no entry -> no allowed
rows, so
// data-2 is skipped entirely (its factory reader is never built).
let calls = std::sync::Mutex::new(Vec::<String>::new());
- let mut factory = as_factory(|f: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- calls.lock().unwrap().push(f.file_name.clone());
- Box::pin(async {
- Ok(Box::new(ArrayReader::new(
- 2,
- vec![Some(vec![1.0, 0.0]), Some(vec![2.0, 0.0])],
- )) as Box<dyn PkVectorReader>)
- })
- });
+ let factory = as_search(
+ |file: &BucketActiveFile,
+ queries: &[&[f32]],
+ metric: VectorSearchMetric,
+ exact_limit: usize,
+ is_excluded: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ calls.lock().unwrap().push(file.file_name.clone());
+ let file_name = file.file_name.clone();
+ let query = queries[0].to_vec();
+ Box::pin(async move {
+ let mut reader =
+ ArrayReader::new(2, vec![Some(vec![1.0, 0.0]),
Some(vec![2.0, 0.0])]);
+ Ok(vec![exact_search(
+ &file_name,
+ &mut reader,
+ &query,
+ metric,
+ exact_limit,
+ is_excluded,
+ )?])
+ })
+ },
+ );
let mut residual: HashMap<String, roaring::RoaringTreemap> =
HashMap::new();
residual.insert("data-1".into(), treemap(&[0, 1]));
let out = bucket_search(
@@ -1086,7 +1211,7 @@ mod tests {
&[],
&[active("data-1", 2), active("data-2", 2)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
5,
@@ -1111,10 +1236,19 @@ mod tests {
// data-1 has an entry but it is empty -> no allowed rows, skipped
without
// reading. Mirrors a file with no residual matches.
let calls = std::sync::Mutex::new(0usize);
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- *calls.lock().unwrap() += 1;
- Box::pin(async { unreachable!("data-1 has an empty allow set and
must not be read") })
- });
+ let factory = as_search(
+ |_: &BucketActiveFile,
+ _: &[&[f32]],
+ _: VectorSearchMetric,
+ _: usize,
+ _: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ *calls.lock().unwrap() += 1;
+ Box::pin(async {
+ unreachable!("data-1 has an empty allow set and must not
be read")
+ })
+ },
+ );
let mut residual: HashMap<String, roaring::RoaringTreemap> =
HashMap::new();
residual.insert("data-1".into(), treemap(&[]));
let out = bucket_search(
@@ -1122,7 +1256,7 @@ mod tests {
&[],
&[active("data-1", 3)],
&HashMap::new(),
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
5,
@@ -1142,18 +1276,11 @@ mod tests {
async fn test_exact_residual_intersects_with_deletion_vector() {
// residual allows {0, 1, 2} but the deletion vector deletes pos0; the
// surviving candidates are the residual-allowed AND not-deleted rows.
- let mut factory = as_factory(|_: &BucketActiveFile| ->
ExactReaderFuture<'_> {
- Box::pin(async {
- Ok(Box::new(ArrayReader::new(
- 2,
- vec![
- Some(vec![1.0, 0.0]),
- Some(vec![2.0, 0.0]),
- Some(vec![3.0, 0.0]),
- ],
- )) as Box<dyn PkVectorReader>)
- })
- });
+ let factory = array_search(vec![
+ Some(vec![1.0, 0.0]),
+ Some(vec![2.0, 0.0]),
+ Some(vec![3.0, 0.0]),
+ ]);
let mut dvs: HashMap<String, Arc<DeletionVector>> = HashMap::new();
let mut bm = RoaringBitmap::new();
bm.insert(0); // pos0 deleted
@@ -1165,7 +1292,7 @@ mod tests {
&[],
&[active("data-1", 3)],
&dvs,
- &mut factory,
+ &factory,
&[0.0, 0.0],
VectorSearchMetric::L2,
5,
@@ -1185,4 +1312,44 @@ mod tests {
vec![1, 2]
);
}
+
+ #[tokio::test]
+ async fn test_non_finite_query_fails_before_opening_any_file() {
+ // A non-finite query element must fail loud before the exact-file
search
+ // closure is ever invoked (validate-before-OPEN). The recording
closure
+ // flips a flag if called; the flag must stay false.
+ let called = std::sync::atomic::AtomicBool::new(false);
+ let factory = as_search(
+ |_: &BucketActiveFile,
+ _: &[&[f32]],
+ _: VectorSearchMetric,
+ _: usize,
+ _: &(dyn Fn(i64) -> bool + Sync)|
+ -> ExactFileSearchFuture<'_> {
+ called.store(true, std::sync::atomic::Ordering::SeqCst);
+ Box::pin(async { unreachable!("closure must not run for a
malformed query") })
+ },
+ );
+ let err = bucket_search(
+ None,
+ &[],
+ &[active("data-1", 2)],
+ &HashMap::new(),
+ &factory,
+ &[f32::NAN, 0.0],
+ VectorSearchMetric::L2,
+ 2,
+ 2,
+ &HashMap::new(),
+ false,
+ None,
+ )
+ .await
+ .unwrap_err();
+ assert!(err.to_string().contains("finite"), "got: {err}");
+ assert!(
+ !called.load(std::sync::atomic::Ordering::SeqCst),
+ "the exact-file search closure must not be invoked for a malformed
query"
+ );
+ }
}
diff --git a/crates/paimon/src/vindex/pkvector/exact.rs
b/crates/paimon/src/vindex/pkvector/exact.rs
index e5d043c6..80333b58 100644
--- a/crates/paimon/src/vindex/pkvector/exact.rs
+++ b/crates/paimon/src/vindex/pkvector/exact.rs
@@ -19,7 +19,10 @@ use std::cmp::Ordering;
use std::collections::BinaryHeap;
use super::data_invalid;
-use super::metric::{java_float_compare, VectorSearchMetric};
+use super::metric::java_float_compare;
+#[cfg(test)]
+use super::metric::VectorSearchMetric;
+#[cfg(test)]
use super::reader::PkVectorReader;
use super::result::PkVectorSearchResult;
@@ -28,7 +31,10 @@ use super::result::PkVectorSearchResult;
/// top therefore evicts the least-wanted candidate. Uses `java_float_compare`
/// for a deterministic total order over f32 that ranks NaN distances as worst
/// (largest), so a NaN is evicted before any finite candidate (no panic).
-struct WorstFirst(PkVectorSearchResult);
+///
+/// Ordered by distance then row position only (a single data file, so the file
+/// name is constant across all its candidates).
+pub(crate) struct WorstFirst(PkVectorSearchResult);
impl PartialEq for WorstFirst {
fn eq(&self, other: &Self) -> bool {
@@ -50,15 +56,74 @@ impl Ord for WorstFirst {
/// True if `candidate` ranks strictly better (BEST_FIRST) than the current
/// worst-on-heap `weakest`: smaller distance, ties broken by smaller position.
-fn is_better_than(candidate: &PkVectorSearchResult, weakest:
&PkVectorSearchResult) -> bool {
+pub(crate) fn is_better_than(
+ candidate: &PkVectorSearchResult,
+ weakest: &PkVectorSearchResult,
+) -> bool {
java_float_compare(candidate.distance, weakest.distance)
.then_with(|| candidate.row_position.cmp(&weakest.row_position))
== Ordering::Less
}
+/// Add `candidate` to a bounded (size `limit`) BEST_FIRST Top-K max-heap over
one
+/// file's candidates: push if under capacity, else replace the current worst
iff
+/// the candidate beats it. `O(log limit)` per call. The single shared push
step
+/// so the whole-file `exact_search` and the streaming per-file search cannot
+/// drift in how they bound their heaps.
+pub(crate) fn push_bounded(
+ heap: &mut BinaryHeap<WorstFirst>,
+ candidate: PkVectorSearchResult,
+ limit: usize,
+) {
+ if heap.len() < limit {
+ heap.push(WorstFirst(candidate));
+ } else if heap
+ .peek()
+ .is_some_and(|worst| is_better_than(&candidate, &worst.0))
+ {
+ heap.pop();
+ heap.push(WorstFirst(candidate));
+ }
+}
+
+/// Drain a bounded Top-K heap into a BEST_FIRST-sorted result list: distance
+/// ASC, then row_position ASC (single file, so data_file_name is constant).
+pub(crate) fn drain_best_first(heap: BinaryHeap<WorstFirst>) ->
Vec<PkVectorSearchResult> {
+ let mut results: Vec<PkVectorSearchResult> = heap.into_iter().map(|w|
w.0).collect();
+ results.sort_by(|a, b| {
+ java_float_compare(a.distance, b.distance).then_with(||
a.row_position.cmp(&b.row_position))
+ });
+ results
+}
+
+/// Validate one query vector against the index dimension: length must match
and
+/// every element must be finite. Shared by the whole-file `exact_search` and
the
+/// table-layer streaming per-file search so both reject the same malformed
+/// queries with the same messages (before any read is performed).
+pub(crate) fn validate_query(query: &[f32], dimension: usize) ->
crate::Result<()> {
+ if query.len() != dimension {
+ return Err(data_invalid(format!(
+ "query vector dimension does not match: index expects {}, got {}",
+ dimension,
+ query.len()
+ )));
+ }
+ if let Some(i) = query.iter().position(|v| !v.is_finite()) {
+ return Err(data_invalid(format!(
+ "query vector element at position {i} must be finite"
+ )));
+ }
+ Ok(())
+}
+
/// Exact Top-K over one sequential physical-row vector source. Mirrors Java
/// `PkVectorExactSearcher.search`. Results are sorted BEST_FIRST: distance
ASC,
/// then row_position ASC (single file, so data_file_name is constant).
+///
+/// No longer the production read path (the table layer streams the vector
column
+/// one Arrow batch at a time via the per-file search closure); kept as a
tested
+/// reference the streaming search must stay byte-identical to.
+#[cfg(test)]
pub(crate) fn exact_search(
data_file_name: &str,
reader: &mut dyn PkVectorReader,
@@ -67,21 +132,10 @@ pub(crate) fn exact_search(
limit: usize,
is_excluded: &dyn Fn(i64) -> bool,
) -> crate::Result<Vec<PkVectorSearchResult>> {
- if query.len() != reader.dimension() {
- return Err(data_invalid(format!(
- "query vector dimension does not match: index expects {}, got {}",
- reader.dimension(),
- query.len()
- )));
- }
+ validate_query(query, reader.dimension())?;
if limit == 0 {
return Err(data_invalid("vector search limit must be positive"));
}
- if let Some(i) = query.iter().position(|v| !v.is_finite()) {
- return Err(data_invalid(format!(
- "query vector element at position {i} must be finite"
- )));
- }
let row_count = reader.row_count();
if row_count < 0 {
return Err(data_invalid(format!(
@@ -101,22 +155,10 @@ pub(crate) fn exact_search(
row_position: position,
distance: metric.compute_distance(query, &reuse),
};
- if heap.len() < limit {
- heap.push(WorstFirst(candidate));
- } else if heap
- .peek()
- .is_some_and(|worst| is_better_than(&candidate, &worst.0))
- {
- heap.pop();
- heap.push(WorstFirst(candidate));
- }
+ push_bounded(&mut heap, candidate, limit);
}
- let mut results: Vec<PkVectorSearchResult> = heap.into_iter().map(|w|
w.0).collect();
- results.sort_by(|a, b| {
- java_float_compare(a.distance, b.distance).then_with(||
a.row_position.cmp(&b.row_position))
- });
- Ok(results)
+ Ok(drain_best_first(heap))
}
#[cfg(test)]
diff --git a/crates/paimon/src/vindex/pkvector/reader.rs
b/crates/paimon/src/vindex/pkvector/reader.rs
index f0f7d0aa..1e0fe8b3 100644
--- a/crates/paimon/src/vindex/pkvector/reader.rs
+++ b/crates/paimon/src/vindex/pkvector/reader.rs
@@ -21,6 +21,10 @@
/// `Send` so a boxed reader can be held across the `.await` points of the
async
/// search path (the returned future is spawned on a `Send` runtime by callers
/// such as the DataFusion integration).
+///
+/// Retained as a test-only reference source for the exact search: production
now
+/// streams the vector column one Arrow batch at a time in the table layer.
+#[cfg(test)]
pub(crate) trait PkVectorReader: Send {
fn dimension(&self) -> usize;