diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 277f04cf5..d18651913 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -18,6 +18,7 @@ use std::collections::{HashMap, HashSet}; const DELETION_VECTORS_ENABLED_OPTION: &str = "deletion-vectors.enabled"; +const DELETION_VECTORS_MERGE_ON_READ_OPTION: &str = "deletion-vectors.merge-on-read"; pub(crate) const QUERY_AUTH_ENABLED_OPTION: &str = "query-auth.enabled"; const DATA_EVOLUTION_ENABLED_OPTION: &str = "data-evolution.enabled"; const GLOBAL_INDEX_ENABLED_OPTION: &str = "global-index.enabled"; @@ -369,6 +370,17 @@ impl<'a> CoreOptions<'a> { .unwrap_or(false) } + /// Whether `deletion-vectors.merge-on-read` is set (default `false`, matching + /// Java `CoreOptions.DELETION_VECTORS_MERGE_ON_READ`). When true, uncompacted + /// (level-0) data is made visible by merging on read; when false, deletion + /// vectors alone determine live rows over the compacted files. + pub fn deletion_vectors_merge_on_read(&self) -> bool { + self.options + .get(DELETION_VECTORS_MERGE_ON_READ_OPTION) + .map(|value| value.eq_ignore_ascii_case("true")) + .unwrap_or(false) + } + /// Whether `query-auth.enabled` is set. /// /// When set, the server enforces a per-user row filter / column masking that this client diff --git a/crates/paimon/src/table/pk_vector_orchestrator.rs b/crates/paimon/src/table/pk_vector_orchestrator.rs index 59ee64cda..6004a98c1 100644 --- a/crates/paimon/src/table/pk_vector_orchestrator.rs +++ b/crates/paimon/src/table/pk_vector_orchestrator.rs @@ -26,6 +26,8 @@ use std::cmp::Ordering; use std::collections::HashMap; use std::sync::Arc; +use roaring::RoaringTreemap; + use crate::deletion_vector::DeletionVector; use crate::spec::BinaryRow; use crate::table::data_file_reader::DataFileReader; @@ -275,6 +277,13 @@ impl PkVectorOrchestrator { /// preserved). The exact-reader factory is split-scoped: it receives the /// current split index and split so a caller can build a reader keyed to the /// specific split/file. `skip_exact_fallback` forwards to `bucket_search`. + /// + /// `residual_by_split`, when present, carries one per-file allow-list of + /// physical row positions per split (indexed parallel to `splits`): only + /// positions listed for a file may survive that bucket's search. A file + /// absent from its split's map (or mapped to an empty set) contributes no + /// candidates. `None` applies no residual filtering. The slice must have the + /// same length as `splits`. #[allow(clippy::too_many_arguments)] #[allow(clippy::type_complexity)] pub(crate) async fn search_candidates( @@ -292,6 +301,7 @@ impl PkVectorOrchestrator { + Send), search_options: &HashMap, skip_exact_fallback: bool, + residual_by_split: Option<&[HashMap]>, ) -> crate::Result> { // Eager input-shape validation (Java checkArgument parity). if limit == 0 { @@ -300,6 +310,13 @@ impl PkVectorOrchestrator { if query.is_empty() { return Err(data_invalid("vector search query must not be empty")); } + if let Some(per_split) = residual_by_split { + if per_split.len() != splits.len() { + return Err(data_invalid( + "residual range map count does not match split count", + )); + } + } // Eager per-bucket search -> tagged candidates. let mut candidates: Vec = Vec::new(); @@ -308,6 +325,7 @@ impl PkVectorOrchestrator { // Wrap the split-scoped factory into bucket_search's per-file signature. let mut bucket_factory = |file: &BucketActiveFile| exact_reader_factory(split_index, split, file); + let residual_ranges = residual_by_split.map(|per_split| &per_split[split_index]); let results = bucket_search( ann_searcher, &split.ann_segments, @@ -319,6 +337,7 @@ impl PkVectorOrchestrator { limit, search_options, skip_exact_fallback, + residual_ranges, )?; for PkVectorSearchResult { data_file_name, @@ -849,6 +868,7 @@ mod e2e_tests { _active_source_files: &HashSet, _dvs: &HashMap>, _opts: &HashMap, + _residual_ranges: Option<&HashMap>, ) -> crate::Result> { Ok(self.hits.clone()) } @@ -875,7 +895,17 @@ mod e2e_tests { // expects; the split index/split are unused here. let mut wrapped = |_: usize, _: &PkVectorSearchSplit, f: &BucketActiveFile| factory(f); let survivors = orch - .search_candidates(splits, query, metric, limit, ann, &mut wrapped, opts, false) + .search_candidates( + splits, + query, + metric, + limit, + ann, + &mut wrapped, + opts, + false, + None, + ) .await?; let indexed_splits = build_indexed_splits(survivors, splits, metric)?; let mut out = Vec::new(); @@ -911,6 +941,7 @@ mod e2e_tests { &mut factory, &opts, false, + None, ) .await .map(|_| ()) @@ -940,6 +971,7 @@ mod e2e_tests { &mut factory, &opts, false, + None, ) .await .map(|_| ()) @@ -1276,6 +1308,7 @@ mod e2e_tests { &mut factory, &opts, false, + None, ) .await .unwrap(); @@ -1289,6 +1322,130 @@ mod e2e_tests { ); } + #[tokio::test] + async fn search_candidates_applies_residual_ranges() { + // Same single bucket / three rows as above, but a per-split residual map + // allows only physical positions {0, 2}. The best recalled hit (pos1, + // d=1) is filtered out; the survivors are the allowed positions in + // best-first order: pos2 (d=4) then pos0 (d=9). This proves the residual + // allow-list is threaded through to the bucket search rather than merely + // stored. + let table_path = "memory:/pkvo_residual"; + let bucket_path = format!("{table_path}/bucket-0"); + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let meta = write_file(&file_io, &bucket_path, "r.mosaic", vec![1, 2, 3]).await; + let split = PkVectorSearchSplit { + data_split: DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(bucket_path) + .with_total_buckets(1) + .with_data_files(vec![meta]) + .build() + .unwrap(), + ann_segments: Vec::new(), + 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 = |_: usize, + _: &PkVectorSearchSplit, + f: &BucketActiveFile| + -> crate::Result> { + assert_eq!(f.file_name, "r.mosaic"); + Ok(Box::new(ArrayReader::new( + 2, + vec![ + Some(vec![3.0, 0.0]), + Some(vec![1.0, 0.0]), + Some(vec![2.0, 0.0]), + ], + ))) + }; + // Allow only positions 0 and 2 for "r.mosaic"; pos1 (the best hit) is + // excluded by the residual. + let mut allowed = RoaringTreemap::new(); + allowed.insert(0); + allowed.insert(2); + let residual_by_split = vec![HashMap::from([("r.mosaic".to_string(), allowed)])]; + let opts = HashMap::new(); + let cands = PkVectorOrchestrator::new(make_reader(file_io, table_path)) + .search_candidates( + &[split], + &[0.0, 0.0], + VectorSearchMetric::L2, + 3, + None, + &mut factory, + &opts, + false, + Some(&residual_by_split), + ) + .await + .unwrap(); + // Best-first among allowed positions: pos2 (d=4) then pos0 (d=9). + assert_eq!( + cands + .iter() + .map(|c| (c.row_position, c.distance)) + .collect::>(), + vec![(2, 4.0), (0, 9.0)] + ); + } + + #[tokio::test] + async fn search_candidates_rejects_residual_length_mismatch() { + // A residual slice whose length differs from the split count is a caller + // bug (the map is indexed by split); it must fail loud rather than panic + // or silently misalign. + let table_path = "memory:/pkvo_residual_mismatch"; + let bucket_path = format!("{table_path}/bucket-0"); + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let meta = write_file(&file_io, &bucket_path, "m.mosaic", vec![1, 2]).await; + let split = PkVectorSearchSplit { + data_split: DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(bucket_path) + .with_total_buckets(1) + .with_data_files(vec![meta]) + .build() + .unwrap(), + ann_segments: Vec::new(), + active_files: vec![active("m.mosaic", 2)], + }; + let mut factory = |_: usize, + _: &PkVectorSearchSplit, + _: &BucketActiveFile| + -> crate::Result> { + unreachable!("length guard must fire before any bucket search") + }; + // Two residual maps for a single split. + let residual_by_split: Vec> = + vec![HashMap::new(), HashMap::new()]; + let opts = HashMap::new(); + let err = PkVectorOrchestrator::new(make_reader(file_io, table_path)) + .search_candidates( + &[split], + &[0.0, 0.0], + VectorSearchMetric::L2, + 3, + None, + &mut factory, + &opts, + false, + Some(&residual_by_split), + ) + .await + .map(|_| ()) + .expect_err("residual length mismatch must fail loud"); + assert!( + format!("{err:?}").contains("does not match split count"), + "got: {err:?}" + ); + } + #[tokio::test] async fn search_candidates_fast_mode_skips_exact_factory() { let table_path = "memory:/pkvo_fast"; @@ -1325,6 +1482,7 @@ mod e2e_tests { &mut factory, &opts, true, + None, ) .await .unwrap(); diff --git a/crates/paimon/src/table/pk_vector_scan.rs b/crates/paimon/src/table/pk_vector_scan.rs index a8219fd31..aebea91b6 100644 --- a/crates/paimon/src/table/pk_vector_scan.rs +++ b/crates/paimon/src/table/pk_vector_scan.rs @@ -24,7 +24,7 @@ use std::collections::{BTreeMap, HashSet}; use crate::spec::{ BinaryRow, DataFileMeta, FileKind, GlobalIndexMeta, IndexManifest, PkVectorSourceFile, - PkVectorSourceMeta, + PkVectorSourceMeta, Predicate, }; use crate::table::pk_vector_orchestrator::PkVectorSearchSplit; use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile}; @@ -204,14 +204,21 @@ pub(crate) struct PkVectorScan<'a> { table: &'a Table, vector_field_id: i32, index_type: String, + filter: Option, } impl<'a> PkVectorScan<'a> { - pub(crate) fn new(table: &'a Table, vector_field_id: i32, index_type: String) -> Self { + pub(crate) fn new( + table: &'a Table, + vector_field_id: i32, + index_type: String, + filter: Option, + ) -> Self { Self { table, vector_field_id, index_type, + filter, } } @@ -226,9 +233,17 @@ impl<'a> PkVectorScan<'a> { // index from it). It also avoids a time-travel mismatch (data from the // travelled snapshot, index from latest) and a TOCTOU where a concurrent // commit lands between two independent resolutions. - let data_splits = self - .table - .new_read_builder() + // + // The residual scalar filter, when set, is pushed into the read builder so + // scan planning drops files whose stats cannot match the predicate, mirroring + // Java `PrimaryKeyVectorScan` applying the filter at scan time. Files that + // survive are still residual-filtered per row downstream; this only avoids + // re-reading files the predicate already excludes. + let mut read_builder = self.table.new_read_builder(); + if let Some(filter) = &self.filter { + read_builder.with_filter(filter.clone()); + } + let data_splits = read_builder .new_scan() .with_scan_all_files() .plan() @@ -613,4 +628,265 @@ mod tests { // Both files are COMPACT + level>0, so both appear as active files. assert_eq!(splits[0].active_files.len(), 2); } + + // ---- Real-table planning tests for filter push-down ---- + // + // Gated off Windows: these fixtures build a table at a `file://` URL derived + // from a temp dir path, which `FileIO` cannot resolve on Windows (see #397). + #[cfg(not(windows))] + mod prune_pushdown_tests { + use super::*; + use crate::catalog::Identifier; + use crate::io::{FileIO, FileIOBuilder}; + use crate::spec::stats::compute_column_stats; + use crate::spec::{ + DataType, Datum, FloatType, IntType, PredicateBuilder, Schema, TableSchema, VectorType, + }; + use crate::table::{CommitMessage, SchemaManager, Table, TableCommit, TableWrite}; + use arrow_array::builder::{FixedSizeListBuilder, Float32Builder}; + use arrow_array::{ArrayRef, Int32Array, RecordBatch}; + use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema}; + use std::sync::Arc; + + /// Vector dimension for the pruning fixtures. + const PRUNE_DIM: usize = 4; + /// The primary-key vector column name. + const PRUNE_VECTOR_COLUMN: &str = "embedding"; + /// vindex index type string; only used to route `PkVectorScan::new`, no index + /// segment is built for these tests. + const PRUNE_INDEX_TYPE: &str = "ivf-flat"; + /// Number of rows written; `id`/`score` values live in `0..PRUNE_ROWS`. + const PRUNE_ROWS: i32 = 4; + /// A predicate literal guaranteed to fall outside the written `id`/`score` + /// range, so file stats cannot match it. + const OUT_OF_RANGE: i32 = 1_000_000; + + /// Schema `(id INT PRIMARY KEY, score INT, embedding VECTOR)`. When + /// `with_deletion_vectors`, enable deletion vectors (merge-on-read left at the + /// default `false`) so a non-PK scalar predicate also stats-prunes; otherwise a + /// plain PK table where only PK-column conjuncts prune. + fn prune_schema(with_deletion_vectors: bool) -> TableSchema { + let mut builder = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("score", DataType::Int(IntType::new())) + .column( + PRUNE_VECTOR_COLUMN, + DataType::Vector( + VectorType::try_new( + true, + PRUNE_DIM as u32, + DataType::Float(FloatType::new()), + ) + .unwrap(), + ), + ) + .primary_key(["id"]) + .option("bucket".to_string(), "1".to_string()); + if with_deletion_vectors { + builder = + builder.option("deletion-vectors.enabled".to_string(), "true".to_string()); + } + TableSchema::new(0, &builder.build().unwrap()) + } + + /// Arrow batch matching the schema: `id` and `score` both equal the physical + /// position (`0..n`), plus a `FixedSizeList` vector column. + fn prune_data_batch(n: usize) -> RecordBatch { + let ids: Vec = (0..n as i32).collect(); + let scores: Vec = (0..n as i32).collect(); + + let element_field = Arc::new(ArrowField::new("element", ArrowDataType::Float32, true)); + let mut vector_builder = + FixedSizeListBuilder::new(Float32Builder::new(), PRUNE_DIM as i32) + .with_field(element_field.clone()); + for i in 0..n { + for d in 0..PRUNE_DIM { + vector_builder.values().append_value((i + d) as f32); + } + vector_builder.append(true); + } + + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("id", ArrowDataType::Int32, false), + ArrowField::new("score", ArrowDataType::Int32, false), + ArrowField::new( + PRUNE_VECTOR_COLUMN, + ArrowDataType::FixedSizeList(element_field, PRUNE_DIM as i32), + true, + ), + ])); + RecordBatch::try_new( + schema, + vec![ + Arc::new(Int32Array::from(ids)) as ArrayRef, + Arc::new(Int32Array::from(scores)) as ArrayRef, + Arc::new(vector_builder.finish()) as ArrayRef, + ], + ) + .unwrap() + } + + async fn prune_open_table(file_io: &FileIO, location: &str) -> Table { + let schema = SchemaManager::new(file_io.clone(), location.to_string()) + .latest() + .await + .expect("failed to list schemas") + .expect("table has no schema"); + Table::new( + file_io.clone(), + Identifier::new("default", "pkvector_prune"), + location.to_string(), + (*schema).clone(), + None, + ) + } + + /// Build a real single-file primary-key table via the public write path, in a + /// fresh temp dir. Persists the schema and writes one data batch, then commits + /// the written data file with real `value_stats` for the `id`/`score` columns. + /// + /// The stats injection mirrors the meta-modification the baseline fixture uses + /// for `level`/`file_source`: the Rust key-value (primary-key) writer records + /// column stats in `key_stats` and leaves `value_stats` empty, but scan-time + /// file pruning reads `value_stats`. Java primary-key writers populate value + /// stats, so committing them here makes the file prunable exactly as it would be + /// in a table written by the Java engine. Returns the temp dir (kept alive by + /// the caller) and the opened table. + async fn build_pruning_test_table( + with_deletion_vectors: bool, + ) -> (tempfile::TempDir, Table) { + let tmp = tempfile::tempdir().expect("create temp dir"); + let location = format!("file://{}", tmp.path().display()); + let file_io = FileIOBuilder::new("file").build().unwrap(); + + for dir in ["schema", "snapshot", "manifest", "index"] { + file_io.mkdirs(&format!("{location}/{dir}")).await.unwrap(); + } + let schema = prune_schema(with_deletion_vectors); + file_io + .new_output(&format!("{location}/schema/schema-{}", schema.id())) + .unwrap() + .write(bytes::Bytes::from(serde_json::to_vec(&schema).unwrap())) + .await + .unwrap(); + + let table = prune_open_table(&file_io, &location).await; + + let batch = prune_data_batch(PRUNE_ROWS as usize); + let mut writer = TableWrite::new(&table, "pkvector-prune".to_string()).unwrap(); + writer.write_arrow_batch(&batch).await.unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + assert_eq!(messages.len(), 1, "single bucket -> one write message"); + let written = &messages[0]; + assert_eq!(written.new_files.len(), 1, "single data file expected"); + let base_meta = written.new_files[0].clone(); + let bucket = written.bucket; + let partition = written.partition.clone(); + + // Real value stats over the `id` (col 0) and `score` (col 1) columns, so a + // predicate outside the written [0, PRUNE_ROWS) range can prune the file. + let int = DataType::Int(IntType::new()); + let value_stats: BinaryTableStats = + compute_column_stats(&batch, &[0, 1], &[int.clone(), int]).unwrap(); + let indexed_meta = DataFileMeta { + value_stats, + value_stats_cols: Some(vec!["id".to_string(), "score".to_string()]), + ..base_meta + }; + + let message = CommitMessage::new(partition, bucket, vec![indexed_meta]); + TableCommit::new(table.clone(), "pkvector-prune".to_string()) + .commit(vec![message]) + .await + .unwrap(); + + (tmp, table) + } + + fn prune_vector_field_id(table: &Table) -> i32 { + table + .schema() + .fields() + .iter() + .find(|f| f.name() == PRUNE_VECTOR_COLUMN) + .expect("vector field present") + .id() + } + + fn prune_equal(table: &Table, column: &str, value: i32) -> Predicate { + PredicateBuilder::new(table.schema().fields()) + .equal(column, Datum::Int(value)) + .unwrap() + } + + #[tokio::test] + async fn plan_prunes_file_when_pk_predicate_excludes_it() { + // Real PK table, one data file with id in [0, PRUNE_ROWS). A predicate + // `id = OUT_OF_RANGE` cannot match the file's id stats, so the scan drops + // the file and plan() returns no splits. Control (no filter) returns one. + let (_tmp, table) = build_pruning_test_table(false).await; + let field_id = prune_vector_field_id(&table); + + let unfiltered = + PkVectorScan::new(&table, field_id, PRUNE_INDEX_TYPE.to_string(), None) + .plan() + .await + .unwrap(); + assert_eq!( + unfiltered.splits.len(), + 1, + "control: file present without a filter" + ); + + let out_of_range = prune_equal(&table, "id", OUT_OF_RANGE); + let filtered = PkVectorScan::new( + &table, + field_id, + PRUNE_INDEX_TYPE.to_string(), + Some(out_of_range), + ) + .plan() + .await + .unwrap(); + assert!( + filtered.splits.is_empty(), + "pk predicate stats-excludes the only file" + ); + } + + #[tokio::test] + async fn plan_prunes_file_on_non_pk_predicate_under_deletion_vectors() { + // Under deletion vectors (merge-on-read off), a non-PK column's stats also + // prune. A `score` predicate outside the written range drops the file. + let (_tmp, table) = build_pruning_test_table(true).await; + let field_id = prune_vector_field_id(&table); + + // Control: without a filter the file is present. + let unfiltered = + PkVectorScan::new(&table, field_id, PRUNE_INDEX_TYPE.to_string(), None) + .plan() + .await + .unwrap(); + assert_eq!( + unfiltered.splits.len(), + 1, + "control: file present without a filter" + ); + + let out_of_range = prune_equal(&table, "score", OUT_OF_RANGE); + let filtered = PkVectorScan::new( + &table, + field_id, + PRUNE_INDEX_TYPE.to_string(), + Some(out_of_range), + ) + .plan() + .await + .unwrap(); + assert!( + filtered.splits.is_empty(), + "non-pk predicate stats-excludes the file under deletion vectors" + ); + } + } } diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 122af0413..facc1f996 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -15,12 +15,14 @@ // specific language governing permissions and limitations // under the License. +use crate::arrow::format::FilePredicates; +use crate::arrow::residual::{filter_record_batch_by_predicates, widen_scan_fields}; use crate::io::FileIO; use crate::lumina::reader::LuminaVectorGlobalIndexReader; use crate::lumina::{is_lumina_index_type, LuminaIndexMeta, LuminaVectorMetric}; use crate::spec::{ - CoreOptions, DataField, FileKind, GlobalIndexSearchMode, IndexFileMeta, IndexManifest, - IndexManifestEntry, ROW_ID_FIELD_NAME, + BigIntType, CoreOptions, DataField, DataType, FileKind, GlobalIndexSearchMode, IndexFileMeta, + IndexManifest, IndexManifestEntry, Predicate, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, }; use crate::table::data_file_reader::DataFileReader; use crate::table::global_index_scanner::{ @@ -38,6 +40,7 @@ use crate::table::pk_vector_position_read::{ }; use crate::table::pk_vector_scan::{PkVectorScan, PkVectorScanPlan}; use crate::table::read_builder::resolve_projected_fields; +use crate::table::source::DataSplit; use crate::table::{ find_field_id_by_name, merge_row_ranges, ArrowRecordBatchStream, RowRange, Table, }; @@ -92,6 +95,7 @@ pub struct VectorSearchBuilder<'a> { limit: Option, options: HashMap, projection: Option>, + filter: Option, } pub struct BatchVectorSearchBuilder<'a> { @@ -111,6 +115,7 @@ impl<'a> VectorSearchBuilder<'a> { limit: None, options: HashMap::new(), projection: None, + filter: None, } } @@ -134,6 +139,26 @@ impl<'a> VectorSearchBuilder<'a> { self } + /// Attach a residual scalar predicate applied *after* vector recall on the + /// primary-key vector path: each recalled candidate file is re-read and only + /// rows satisfying `filter` survive, folded into the search so best-first + /// order and Top-K still hold. Mirrors Java `PrimaryKeyVectorRead`'s + /// residual-filter support. Only the primary-key vector path consumes it, and + /// only when the table exposes physical rows directly (deletion vectors + /// enabled without merge-on-read); otherwise the query fails loud. A query + /// that does not resolve to the primary-key vector path (no PK-vector index, + /// or a non-PK-vector column) also fails loud rather than silently ignoring + /// the filter. + /// + /// The whole predicate is both pushed into the scan — where it prunes whole + /// data files by their column stats — and applied per row as a residual over + /// the surviving files, so results stay exact. Sub-file row-range narrowing is + /// not performed; a surviving file is re-read in full for the residual. + pub fn with_filter(&mut self, filter: Predicate) -> &mut Self { + self.filter = Some(filter); + self + } + /// Restrict the columns materialized by [`execute_read`](Self::execute_read) /// to `cols` (plus the always-appended `_PKEY_VECTOR_SCORE`). Without this /// call `execute_read` materializes every user table column. Only affects @@ -191,6 +216,19 @@ impl<'a> VectorSearchBuilder<'a> { } } + // The data-evolution (global-index) fall-through path cannot honor a + // residual filter — it never reads physical rows. Rather than silently + // drop the predicate and return unfiltered results, fail loud when a + // filter is set on a query that does not resolve to the primary-key + // vector path. + if self.filter.is_some() { + return Err(crate::Error::DataInvalid { + message: "vector search filter is only supported on the primary-key vector path" + .to_string(), + source: None, + }); + } + let mut batch_builder = BatchVectorSearchBuilder::new(self.table); let mut results = batch_builder .with_vector_column(vector_column) @@ -284,11 +322,26 @@ impl<'a> VectorSearchBuilder<'a> { query_vector: &[f32], limit: usize, ) -> crate::Result<(Vec, PkVectorScanPlan, VectorSearchMetric)> { - // Residual-filter guard: PK vector search accepts partition filters only. - // This builder exposes no data-predicate setter, so there is nothing to - // reject here; the guard mirrors Java `checkArgument(filter == null)` and, - // if a filter setter is ever added, it must error rather than be ignored. - + // Residual pre-filter guard, mirroring Java `PrimaryKeyVectorScan`. A data + // predicate set via `with_filter` is applied post-recall by re-reading + // each candidate file's physical rows (see below). That physical-position + // filtering only agrees with the bucket search when the table exposes + // physical rows directly: deletion vectors enabled and merge-on-read + // disabled. Under merge-on-read (or without deletion vectors) a read + // merges multiple key versions, so a scalar filter could retain a stale + // version whose live version does not match — a silent wrong-read. Reject + // such queries rather than answer them incorrectly. No filter → nothing to + // guard, so the search-only and read paths are unaffected. + let physical_row_read = + core.deletion_vectors_enabled() && !core.deletion_vectors_merge_on_read(); + if self.filter.is_some() && !physical_row_read { + return Err(crate::Error::DataInvalid { + message: + "primary-key vector pre-filter requires deletion vectors without merge-on-read" + .to_string(), + source: None, + }); + } // `primary_key_vector_distance_metric` returns a validated name; re-parse // into the enum for the numeric semantics. let metric = VectorSearchMetric::parse(&core.primary_key_vector_distance_metric(pk_col)?)?; @@ -315,7 +368,7 @@ 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) + let plan = PkVectorScan::new(self.table, field_id, index_type, self.filter.clone()) .plan() .await?; if plan.splits.is_empty() { @@ -365,12 +418,61 @@ impl<'a> VectorSearchBuilder<'a> { }); let ann_searcher = VindexAnnSearcher::new(field_name, scorer); + // Residual (post-recall) filtering: for each candidate file, re-read its + // physical rows and keep the positions whose rows satisfy the filter. The + // per-split allow-list is threaded into the bucket search so the residual + // folds into recall (best-first order and Top-K are preserved). Built only + // when a filter is set; otherwise `None` leaves the search unfiltered. The + // residual reader projects the predicate columns plus `_ROW_ID` (used to + // recover file-local physical positions) and carries no pushdown, matching + // `residual_positions_by_file`. Computed before the exact-reader preload so + // the preload can skip files the residual allow-list leaves empty. + let residual_by_split: Option>> = match &self.filter { + Some(filter) => { + let file_predicates = FilePredicates { + predicates: vec![filter.clone()], + file_fields: self.table.schema().fields().to_vec(), + }; + let row_id_field = DataField::new( + ROW_ID_FIELD_ID, + ROW_ID_FIELD_NAME.to_string(), + DataType::BigInt(BigIntType::new()), + ); + let residual_read_type = + widen_scan_fields(std::slice::from_ref(&row_id_field), Some(&file_predicates)); + let residual_reader = DataFileReader::new( + self.table.file_io().clone(), + self.table.schema_manager().clone(), + self.table.schema().id(), + self.table.schema().fields().to_vec(), + residual_read_type, + Vec::new(), + ); + let mut per_split = Vec::with_capacity(plan.splits.len()); + for split in &plan.splits { + per_split.push( + residual_positions_by_file( + &residual_reader, + &split.data_split, + &split.active_files, + &file_predicates, + ) + .await?, + ); + } + Some(per_split) + } + None => None, + }; + // Exact-fallback readers, keyed by (split_index, file_name). In FAST mode // the kernel never invokes the factory, so skip the in-memory column read // entirely. Otherwise preload only the *uncovered* active files: files an // ANN segment already covers never reach the exact fallback, so reading // their vector column here would be wasted IO/memory. Mirrors Java, which - // creates a `PkVectorReader` lazily only for uncovered files. + // creates a `PkVectorReader` lazily only for uncovered files. When a + // residual filter leaves a file's allow-list empty (or absent) the bucket + // search skips it, so its reader is not preloaded either. let mut exact_readers: HashMap<(usize, String), Box> = HashMap::new(); if !skip_exact_fallback { for (split_index, split) in plan.splits.iter().enumerate() { @@ -384,6 +486,13 @@ impl<'a> VectorSearchBuilder<'a> { if covered.contains(&active.file_name) { continue; } + if !should_preload_exact_reader( + residual_by_split.as_deref(), + split_index, + &active.file_name, + ) { + continue; + } let r = factory.create(active).await?; exact_readers.insert((split_index, active.file_name.clone()), r); } @@ -411,6 +520,7 @@ impl<'a> VectorSearchBuilder<'a> { &mut factory, &search_options, skip_exact_fallback, + residual_by_split.as_deref(), ) .await?; @@ -912,6 +1022,114 @@ fn is_vector_global_index_file(index_file: &IndexFileMeta) -> bool { VectorIndexBackend::from_index_type(&index_file.index_type).is_some() } +/// Whether the exact-fallback reader for `file_name` in split `split_index` +/// should be preloaded. With a residual filter, a file absent from the split's +/// allow-list or with an empty allow-list has no candidate rows, so the bucket +/// search skips it and preloading its vector column would be wasted IO. +fn should_preload_exact_reader( + residual_by_split: Option<&[HashMap]>, + split_index: usize, + file_name: &str, +) -> bool { + match residual_by_split { + None => true, + Some(per_split) => per_split + .get(split_index) + .and_then(|m| m.get(file_name)) + .is_some_and(|allowed| !allowed.is_empty()), + } +} + +/// Compute, per data file in `split`, the set of physical row positions whose +/// rows satisfy the residual predicate. Mirrors the row-collecting half of Java +/// `PrimaryKeyVectorRead`'s `executeFilter`: because +/// [`DataFileReader::read_single_file_stream`] rejects projecting `_ROW_ID` +/// alongside a row-filtering predicate (the residual filter would drop rows +/// before `_ROW_ID` is assigned positionally, desyncing it), the predicate is +/// NOT pushed down. Instead `reader` projects the residual columns together with +/// `_ROW_ID` and carries no pushdown predicate; the residual is applied here at +/// the Arrow level, after `_ROW_ID` is materialized, and each surviving row's +/// `_ROW_ID - first_row_id` is the file-local physical position. +/// +/// Every *active* data file in the split gets an entry, possibly empty. The +/// bucket search treats an absent entry and an empty entry identically (the file +/// contributes no candidates), so the empty entries only make the map cover every +/// active file. Non-active files (e.g. level-0 files the bucket search excludes) +/// are skipped entirely: they are never searched, so re-reading them would be +/// wasted IO and their possibly-absent `first_row_id` must not fail an otherwise +/// valid query. +/// +/// `reader` must project `_ROW_ID` and be predicate-free; `residual.file_fields` +/// are the fields the residual leaf indices point into (resolved by name against +/// each emitted batch). A data file without `first_row_id` fails loud, matching +/// the position-read guard. +async fn residual_positions_by_file( + reader: &DataFileReader, + split: &DataSplit, + active_files: &[BucketActiveFile], + residual: &FilePredicates, +) -> crate::Result> { + let scan_fields = reader.read_type().to_vec(); + let active_names: HashSet<&str> = active_files.iter().map(|f| f.file_name.as_str()).collect(); + let mut out: HashMap = HashMap::new(); + for file_meta in split.data_files() { + // Only files the bucket search actually recalls from need residual + // positions; skip everything else so a non-active file cannot trigger the + // `first_row_id` guard below or incur a wasted read. + if !active_names.contains(file_meta.file_name.as_str()) { + continue; + } + let first_row_id = file_meta + .first_row_id + .ok_or_else(|| crate::Error::DataInvalid { + message: format!( + "residual position read requires data file '{}' to have first_row_id", + file_meta.file_name + ), + source: None, + })?; + let data_fields = reader.derive_data_fields(file_meta).await?; + let mut stream = + reader.read_single_file_stream(split, file_meta.clone(), data_fields, None, None)?; + // Register the file up front so a file whose rows all fail the residual + // still appears in the map (empty set). + let positions = out.entry(file_meta.file_name.clone()).or_default(); + while let Some(batch) = stream.try_next().await? { + let filtered = filter_record_batch_by_predicates(batch, residual, &scan_fields)?; + if filtered.num_rows() == 0 { + continue; + } + let row_id_idx = filtered.schema().index_of(ROW_ID_FIELD_NAME).map_err(|_| { + crate::Error::DataInvalid { + message: "residual position read batch is missing the _ROW_ID column" + .to_string(), + source: None, + } + })?; + let row_ids = filtered + .column(row_id_idx) + .as_any() + .downcast_ref::() + .ok_or_else(|| crate::Error::DataInvalid { + message: "residual position read _ROW_ID column is not Int64".to_string(), + source: None, + })?; + for i in 0..row_ids.len() { + let position = row_ids.value(i) - first_row_id; + let position = u64::try_from(position).map_err(|_| crate::Error::DataInvalid { + message: format!( + "residual position {position} is negative for data file '{}'", + file_meta.file_name + ), + source: None, + })?; + positions.insert(position); + } + } + } + Ok(out) +} + /// Preload every ANN segment's bytes into a map keyed by the resolved (globally /// unique) segment path. The scorer closure reads from this map so the vindex /// reader is driven from memory without per-search IO. @@ -1954,8 +2172,8 @@ mod tests { use crate::lumina::{LEGACY_LUMINA_VECTOR_ANN_IDENTIFIER, LUMINA_IDENTIFIER}; use crate::spec::stats::BinaryTableStats; use crate::spec::{ - ArrayType, BinaryRow, DataFileMeta, DataType, FloatType, GlobalIndexMeta, IndexFileMeta, - IndexManifestEntry, IntType, Schema, TableSchema, + ArrayType, BinaryRow, DataFileMeta, DataType, Datum, FloatType, GlobalIndexMeta, + IndexFileMeta, IndexManifestEntry, IntType, PredicateBuilder, Schema, TableSchema, }; use crate::table::source::DataSplitBuilder; use crate::vindex::IVF_FLAT_IDENTIFIER; @@ -2014,6 +2232,26 @@ mod tests { assert_eq!(find_field_id_by_name(&fields, "nonexistent"), None); } + #[test] + fn should_preload_skips_empty_or_absent_residual_files() { + use roaring::RoaringTreemap; + use std::collections::HashMap; + // No filter -> always preload. + assert!(should_preload_exact_reader(None, 0, "f0")); + // Filter present: file with a non-empty allow-list -> preload. + let mut m0: HashMap = HashMap::new(); + m0.insert("f0".to_string(), RoaringTreemap::from_iter([0u64])); + let per_split = vec![m0]; + assert!(should_preload_exact_reader(Some(&per_split), 0, "f0")); + // Filter present: file absent -> skip. + assert!(!should_preload_exact_reader(Some(&per_split), 0, "missing")); + // Filter present: file with empty allow-list -> skip. + let mut m1: HashMap = HashMap::new(); + m1.insert("f1".to_string(), RoaringTreemap::new()); + let per_split2 = vec![m1]; + assert!(!should_preload_exact_reader(Some(&per_split2), 0, "f1")); + } + #[test] fn test_raw_vector_score_matches_java_metric_semantics() { let l2 = compute_raw_vector_score(&[1.0, 2.0], &[1.0, 4.0], RawVectorMetric::L2); @@ -2817,6 +3055,147 @@ mod tests { } } + /// `id > threshold` built against the table's user fields (leaf index resolves + /// against `table.schema().fields()`). + fn id_gt_filter(table: &Table, threshold: i32) -> Predicate { + PredicateBuilder::new(table.schema().fields()) + .greater_than("id", Datum::Int(threshold)) + .unwrap() + } + + #[tokio::test] + async fn pk_branch_filter_without_deletion_vectors_fails_loud() { + // A residual filter on a PK-vector table that does NOT enable deletion + // vectors must be rejected (merge-on-read semantics would make physical + // -position filtering unsound). Mirrors Java `PrimaryKeyVectorScan`. + let table = pk_vector_table(&[ + ("pk-vector.index.columns", "embedding"), + ("fields.embedding.pk-vector.index.type", IVF_FLAT_IDENTIFIER), + ("fields.embedding.pk-vector.distance.metric", "l2"), + ]); + let filter = id_gt_filter(&table, 2); + let err = table + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0]) + .with_limit(5) + .with_filter(filter) + .execute_scored() + .await + .map(|_| ()) + .expect_err("filter without deletion vectors must fail loud"); + assert!( + matches!(err, crate::Error::DataInvalid { ref message, .. } + if message.contains("deletion vectors without merge-on-read")), + "unexpected error: {err:?}" + ); + } + + #[tokio::test] + async fn execute_read_filter_without_deletion_vectors_fails_loud() { + let table = pk_vector_table(&[ + ("pk-vector.index.columns", "embedding"), + ("fields.embedding.pk-vector.index.type", IVF_FLAT_IDENTIFIER), + ("fields.embedding.pk-vector.distance.metric", "l2"), + ]); + let filter = id_gt_filter(&table, 2); + let err = table + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0]) + .with_limit(5) + .with_filter(filter) + .execute_read() + .await + .map(|_| ()) + .expect_err("read filter without deletion vectors must fail loud"); + assert!( + matches!(err, crate::Error::DataInvalid { ref message, .. } + if message.contains("deletion vectors without merge-on-read")), + "unexpected error: {err:?}" + ); + } + + #[tokio::test] + async fn execute_scored_filter_on_non_pk_vector_path_fails_loud() { + // No PK-vector index configured, so `execute_scored` would fall through to + // the data-evolution path, which never consumes the filter. Silently + // returning unfiltered rows is a wrong-read; the query must fail loud + // instead. + let table = pk_vector_table(&[]); + let filter = id_gt_filter(&table, 2); + let err = table + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0]) + .with_limit(5) + .with_filter(filter) + .execute_scored() + .await + .map(|_| ()) + .expect_err("filter on the non-PK-vector path must fail loud"); + assert!( + matches!(err, crate::Error::DataInvalid { ref message, .. } + if message.contains("only supported on the primary-key vector path")), + "unexpected error: {err:?}" + ); + } + + #[tokio::test] + async fn pk_branch_filter_with_merge_on_read_fails_loud() { + // Deletion vectors enabled BUT merge-on-read on: still rejected, because a + // merge-on-read scan can surface stale key versions that a physical-row + // filter cannot reconcile. + let table = pk_vector_table(&[ + ("pk-vector.index.columns", "embedding"), + ("fields.embedding.pk-vector.index.type", IVF_FLAT_IDENTIFIER), + ("fields.embedding.pk-vector.distance.metric", "l2"), + ("deletion-vectors.enabled", "true"), + ("deletion-vectors.merge-on-read", "true"), + ]); + let filter = id_gt_filter(&table, 2); + let err = table + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0]) + .with_limit(5) + .with_filter(filter) + .execute_scored() + .await + .map(|_| ()) + .expect_err("merge-on-read filter must fail loud"); + assert!( + matches!(err, crate::Error::DataInvalid { ref message, .. } + if message.contains("deletion vectors without merge-on-read")), + "unexpected error: {err:?}" + ); + } + + #[tokio::test] + async fn pk_branch_filter_with_deletion_vectors_passes_guard() { + // Deletion vectors enabled, merge-on-read off (default): the residual guard + // passes. With no snapshot the plan is empty, so the (guarded) filter path + // simply yields an empty result rather than erroring — proving the guard + // admits a legal filtered query. + let table = pk_vector_table(&[ + ("pk-vector.index.columns", "embedding"), + ("fields.embedding.pk-vector.index.type", IVF_FLAT_IDENTIFIER), + ("fields.embedding.pk-vector.distance.metric", "l2"), + ("deletion-vectors.enabled", "true"), + ]); + let filter = id_gt_filter(&table, 2); + let result = table + .new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector(vec![1.0]) + .with_limit(5) + .with_filter(filter) + .execute_scored() + .await + .expect("guarded filter query must be admitted"); + assert!(result.is_empty()); + } + fn make_lumina_entry( file_name: &str, index_type: &str, @@ -3114,3 +3493,324 @@ mod tests { assert_eq!(names, vec!["id"]); } } + +/// Tests for [`residual_positions_by_file`]: the residual predicate is applied at +/// the Arrow level (no pushdown) against the predicate columns plus `_ROW_ID`, +/// and each surviving row's `_ROW_ID` is converted back to a file-local physical +/// position. +#[cfg(test)] +mod residual_positions_tests { + use super::*; + use crate::arrow::build_target_arrow_schema; + use crate::arrow::format::FilePredicates; + use crate::io::FileIOBuilder; + use crate::spec::stats::BinaryTableStats; + use crate::spec::{ + BigIntType, BinaryRow, DataField, DataFileMeta, DataType, Datum, IntType, PredicateBuilder, + ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, + }; + use crate::table::data_file_reader::DataFileReader; + use crate::table::schema_manager::SchemaManager; + use crate::table::source::{DataSplit, DataSplitBuilder}; + use arrow_array::{Int32Array, RecordBatch}; + use bytes::Bytes; + use paimon_mosaic_core::spec::COMPRESSION_NONE; + use paimon_mosaic_core::writer::{MosaicWriter, OutputFile, WriterOptions}; + use std::io; + use std::sync::Arc; + + struct MemOutputFile { + data: Vec, + } + + impl OutputFile for MemOutputFile { + fn write(&mut self, data: &[u8]) -> io::Result<()> { + self.data.extend_from_slice(data); + Ok(()) + } + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + fn pos(&self) -> u64 { + self.data.len() as u64 + } + } + + fn id_field() -> DataField { + DataField::new(0, "id".to_string(), DataType::Int(IntType::new())) + } + + fn row_id_field() -> DataField { + DataField::new( + ROW_ID_FIELD_ID, + ROW_ID_FIELD_NAME.to_string(), + DataType::BigInt(BigIntType::new()), + ) + } + + fn id_batch(ids: Vec) -> RecordBatch { + let schema = build_target_arrow_schema(&[id_field()]).unwrap(); + RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(ids))]).unwrap() + } + + fn write_mosaic(batch: &RecordBatch) -> Bytes { + let mut writer = MosaicWriter::new( + MemOutputFile { data: Vec::new() }, + batch.schema().as_ref(), + WriterOptions { + compression: COMPRESSION_NONE, + num_buckets: 2, + row_group_max_size: u64::MAX, + ..Default::default() + }, + ) + .unwrap(); + writer.write_batch(batch).unwrap(); + writer.close().unwrap(); + Bytes::from(writer.output().data.to_vec()) + } + + fn data_file( + file_name: &str, + file_size: i64, + row_count: i64, + first_row_id: Option, + ) -> DataFileMeta { + DataFileMeta { + file_name: file_name.to_string(), + file_size, + row_count, + min_key: Vec::new(), + max_key: Vec::new(), + key_stats: BinaryTableStats::empty(), + value_stats: BinaryTableStats::empty(), + min_sequence_number: 0, + max_sequence_number: 0, + schema_id: 1, + level: 0, + extra_files: Vec::new(), + creation_time: None, + delete_row_count: None, + embedded_index: None, + file_source: None, + value_stats_cols: None, + external_path: None, + first_row_id, + write_cols: None, + } + } + + /// Build a predicate-free reader (read_type = `id` + `_ROW_ID`) over a split + /// containing `files` (each `(name, ids, first_row_id)`), written as Mosaic + /// data files in the same bucket. The returned active-file list covers every + /// file (all files active). + async fn build_reader_and_split( + table_path: &str, + files: &[(&str, Vec, i64)], + ) -> (DataFileReader, DataSplit, Vec) { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let bucket_path = format!("{table_path}/bucket-0"); + let mut metas = Vec::new(); + let mut active_files = Vec::new(); + for (name, ids, first_row_id) in files { + let data = write_mosaic(&id_batch(ids.clone())); + file_io + .new_output(&format!("{bucket_path}/{name}")) + .unwrap() + .write(data.clone()) + .await + .unwrap(); + metas.push(data_file( + name, + data.len() as i64, + ids.len() as i64, + Some(*first_row_id), + )); + active_files.push(BucketActiveFile { + file_name: name.to_string(), + row_count: ids.len() as i64, + }); + } + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(bucket_path) + .with_total_buckets(1) + .with_data_files(metas) + .build() + .unwrap(); + let reader = DataFileReader::new( + file_io.clone(), + SchemaManager::new(file_io, table_path.to_string()), + 1, + vec![id_field()], + vec![id_field(), row_id_field()], + Vec::new(), + ); + (reader, split, active_files) + } + + /// `id > threshold`, with `file_fields` = `[id]` so the leaf index resolves. + fn residual_id_gt(threshold: i32) -> FilePredicates { + let pred = PredicateBuilder::new(&[id_field()]) + .greater_than("id", Datum::Int(threshold)) + .unwrap(); + FilePredicates { + predicates: vec![pred], + file_fields: vec![id_field()], + } + } + + fn sorted(t: &roaring::RoaringTreemap) -> Vec { + t.iter().collect() + } + + #[tokio::test] + async fn test_residual_selects_matching_positions() { + // ids [1,2,3,4,5] at first_row_id 0; id > 2 -> ids 3,4,5 -> positions 2,3,4. + let (reader, split, active) = build_reader_and_split( + "memory:/rpf_basic", + &[("part-0.mosaic", vec![1, 2, 3, 4, 5], 0)], + ) + .await; + let map = residual_positions_by_file(&reader, &split, &active, &residual_id_gt(2)) + .await + .unwrap(); + assert_eq!(sorted(&map["part-0.mosaic"]), vec![2, 3, 4]); + } + + #[tokio::test] + async fn test_residual_matches_none_yields_empty_entry() { + // id > 100 matches nothing; the file still gets a (present, empty) entry. + let (reader, split, active) = + build_reader_and_split("memory:/rpf_none", &[("part-0.mosaic", vec![1, 2, 3], 0)]) + .await; + let map = residual_positions_by_file(&reader, &split, &active, &residual_id_gt(100)) + .await + .unwrap(); + assert!(map.contains_key("part-0.mosaic")); + assert!(map["part-0.mosaic"].is_empty()); + } + + #[tokio::test] + async fn test_residual_matches_all_yields_full_set() { + let (reader, split, active) = + build_reader_and_split("memory:/rpf_all", &[("part-0.mosaic", vec![1, 2, 3], 0)]).await; + let map = residual_positions_by_file(&reader, &split, &active, &residual_id_gt(0)) + .await + .unwrap(); + assert_eq!(sorted(&map["part-0.mosaic"]), vec![0, 1, 2]); + } + + #[tokio::test] + async fn test_residual_positions_are_file_local_across_files() { + // Two files with distinct first_row_id; positions must be 0-based within + // each file, not global. id > 3 keeps ids 4,5 in both -> positions {3,4}. + let (reader, split, active) = build_reader_and_split( + "memory:/rpf_multi", + &[ + ("part-0.mosaic", vec![1, 2, 3, 4, 5], 0), + ("part-1.mosaic", vec![1, 2, 3, 4, 5], 100), + ], + ) + .await; + let map = residual_positions_by_file(&reader, &split, &active, &residual_id_gt(3)) + .await + .unwrap(); + assert_eq!(sorted(&map["part-0.mosaic"]), vec![3, 4]); + assert_eq!(sorted(&map["part-1.mosaic"]), vec![3, 4]); + } + + #[tokio::test] + async fn test_non_active_files_are_skipped() { + // Two files in the split, but only `part-0.mosaic` is active. The bucket + // search never recalls from `part-1.mosaic` (level-0 / non-active), so it + // must not appear in the residual map — and even though it lacks a + // `first_row_id`, the query still succeeds because non-active files are + // skipped before the guard. + let (reader, split, mut active) = build_reader_and_split( + "memory:/rpf_nonactive", + &[("part-0.mosaic", vec![1, 2, 3, 4, 5], 0)], + ) + .await; + // Append a non-active file (missing first_row_id) directly to the split's + // data files, but leave it out of the active list. + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let bucket_path = "memory:/rpf_nonactive/bucket-0"; + let data = write_mosaic(&id_batch(vec![9, 9, 9])); + file_io + .new_output(&format!("{bucket_path}/part-1.mosaic")) + .unwrap() + .write(data.clone()) + .await + .unwrap(); + let mut metas = split.data_files().to_vec(); + metas.push(data_file("part-1.mosaic", data.len() as i64, 3, None)); + // `active` already lists only part-0.mosaic; keep it that way. + let _ = &mut active; + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(bucket_path.to_string()) + .with_total_buckets(1) + .with_data_files(metas) + .build() + .unwrap(); + let map = residual_positions_by_file(&reader, &split, &active, &residual_id_gt(2)) + .await + .unwrap(); + assert_eq!(sorted(&map["part-0.mosaic"]), vec![2, 3, 4]); + assert!( + !map.contains_key("part-1.mosaic"), + "non-active file must be skipped" + ); + } + + #[tokio::test] + async fn test_missing_first_row_id_is_error() { + let (reader, split, active) = build_reader_and_split_no_first_row_id().await; + let err = residual_positions_by_file(&reader, &split, &active, &residual_id_gt(0)) + .await + .expect_err("missing first_row_id must error"); + assert!(format!("{err:?}").contains("first_row_id"), "got: {err:?}"); + } + + async fn build_reader_and_split_no_first_row_id( + ) -> (DataFileReader, DataSplit, Vec) { + let table_path = "memory:/rpf_nofrid"; + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let bucket_path = format!("{table_path}/bucket-0"); + let data = write_mosaic(&id_batch(vec![1, 2, 3])); + file_io + .new_output(&format!("{bucket_path}/part-0.mosaic")) + .unwrap() + .write(data.clone()) + .await + .unwrap(); + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(bucket_path) + .with_total_buckets(1) + .with_data_files(vec![data_file("part-0.mosaic", data.len() as i64, 3, None)]) + .build() + .unwrap(); + let reader = DataFileReader::new( + file_io.clone(), + SchemaManager::new(file_io, table_path.to_string()), + 1, + vec![id_field()], + vec![id_field(), row_id_field()], + Vec::new(), + ); + // The lone file is active, so the `first_row_id` guard applies to it. + let active = vec![BucketActiveFile { + file_name: "part-0.mosaic".to_string(), + row_count: 3, + }]; + (reader, split, active) + } +} diff --git a/crates/paimon/src/vindex/pkvector/ann.rs b/crates/paimon/src/vindex/pkvector/ann.rs index aad463f4b..627c052eb 100644 --- a/crates/paimon/src/vindex/pkvector/ann.rs +++ b/crates/paimon/src/vindex/pkvector/ann.rs @@ -35,12 +35,21 @@ use crate::vector_search::VectorSearch; /// longer readable in this snapshot). Deletion vectors are applied only to active /// sources. /// -/// Returns `None` only when every source file is active AND no deletion vector is -/// relevant — nothing to mask. Otherwise returns the masked live ids. +/// `residual_ranges` (when `Some`) restricts each source file to the physical row +/// positions allowed by a residual predicate on the data columns: `key = file +/// name`, `value = allowed physical positions`. A file with no entry (or an empty +/// entry) has no allowed rows and contributes nothing. When `residual_ranges` is +/// `Some`, a mask is always required (the residual can only narrow the live set), +/// so the result is always `Some`. Mirrors Java `rowRangesByFile`. +/// +/// Returns `None` only when there is no residual, every source file is active, AND +/// no deletion vector is relevant — nothing to mask. Otherwise returns the masked +/// live ids. pub(crate) fn build_live_row_ids( source_files: &[PkVectorSourceFile], active_source_files: &HashSet, deletion_vectors: &HashMap>, + residual_ranges: Option<&HashMap>, ) -> crate::Result> { let all_active = source_files .iter() @@ -48,7 +57,7 @@ pub(crate) fn build_live_row_ids( let has_relevant_dv = source_files .iter() .any(|f| deletion_vectors.contains_key(f.file_name())); - if all_active && !has_relevant_dv { + if residual_ranges.is_none() && all_active && !has_relevant_dv { return Ok(None); } @@ -63,7 +72,32 @@ pub(crate) fn build_live_row_ids( .ok_or_else(|| data_invalid("vector source row counts overflow u64"))?; let active = active_source_files.contains(source_file.file_name()); if active && row_count > 0 { - live.insert_range(file_offset..end); + match residual_ranges { + // No residual: the whole active file range is live. + None => { + live.insert_range(file_offset..end); + } + // Residual present: only allowed physical positions of this file + // become live, mapped into global ordinal space (position + + // file_offset). A missing/empty entry allows no rows. + Some(ranges) => { + if let Some(allowed) = ranges.get(source_file.file_name()) { + for position in allowed.iter() { + if position >= row_count { + return Err(data_invalid(format!( + "residual position {position} is out of range for source file {} ({} rows)", + source_file.file_name(), + row_count + ))); + } + let global = file_offset.checked_add(position).ok_or_else(|| { + data_invalid("vector residual position overflows u64") + })?; + live.insert(global); + } + } + } + } } if active { if let Some(dv) = deletion_vectors.get(source_file.file_name()) { @@ -91,6 +125,7 @@ pub(crate) fn map_ann_results( source_meta: &PkVectorSourceMeta, active_source_files: &HashSet, deletion_vectors: &HashMap>, + residual_ranges: Option<&HashMap>, metric: VectorSearchMetric, ) -> crate::Result> { let mut results = Vec::with_capacity(scored.len()); @@ -103,15 +138,23 @@ pub(crate) fn map_ann_results( "ANN segment returned inactive source {data_file_name}" ))); } + let pos = u64::try_from(row_position) + .map_err(|_| data_invalid("resolved row position must not be negative"))?; if let Some(dv) = deletion_vectors.get(&data_file_name) { - let pos = u64::try_from(row_position) - .map_err(|_| data_invalid("resolved row position must not be negative"))?; if dv.is_deleted(pos) { return Err(data_invalid(format!( "ANN segment returned snapshot-deleted row position {row_position} in {data_file_name}" ))); } } + if let Some(ranges) = residual_ranges { + let allowed = ranges.get(&data_file_name).is_some_and(|r| r.contains(pos)); + if !allowed { + return Err(data_invalid(format!( + "ANN segment returned row position {row_position} in {data_file_name} outside the residual pre-filter" + ))); + } + } results.push(PkVectorSearchResult { data_file_name, row_position, @@ -143,6 +186,7 @@ pub(crate) trait PkVectorAnnSearcher: Send + Sync { active_source_files: &HashSet, deletion_vectors: &HashMap>, search_options: &HashMap, + residual_ranges: Option<&HashMap>, ) -> crate::Result>; } @@ -185,6 +229,7 @@ impl PkVectorAnnSearcher for VindexAnnSearcher { active_source_files: &HashSet, deletion_vectors: &HashMap>, search_options: &HashMap, + residual_ranges: Option<&HashMap>, ) -> crate::Result> { if limit == 0 { return Err(data_invalid("vector search limit must be positive")); @@ -192,8 +237,12 @@ impl PkVectorAnnSearcher for VindexAnnSearcher { let source_files = segment.source_meta.source_files(); let mut search = VectorSearch::new(query.to_vec(), limit, self.field_name.clone())? .with_options(search_options.clone()); - if let Some(live) = build_live_row_ids(source_files, active_source_files, deletion_vectors)? - { + if let Some(live) = build_live_row_ids( + source_files, + active_source_files, + deletion_vectors, + residual_ranges, + )? { search = search.with_include_row_ids(live); } let scored = match (self.scorer)(segment, &search)? { @@ -206,6 +255,7 @@ impl PkVectorAnnSearcher for VindexAnnSearcher { &segment.source_meta, active_source_files, deletion_vectors, + residual_ranges, metric, ) } @@ -241,13 +291,15 @@ mod tests { let files = [PkVectorSourceFile::new("f0".into(), 3).unwrap()]; let active = active_set(&["f0"]); // All active + empty map -> None. - assert!(build_live_row_ids(&files, &active, &HashMap::new()) + assert!(build_live_row_ids(&files, &active, &HashMap::new(), None) .unwrap() .is_none()); // All active + non-empty map but no matching file name -> None. let mut dvs = HashMap::new(); dvs.insert("other".to_string(), dv(&[0])); - assert!(build_live_row_ids(&files, &active, &dvs).unwrap().is_none()); + assert!(build_live_row_ids(&files, &active, &dvs, None) + .unwrap() + .is_none()); } #[test] @@ -258,7 +310,7 @@ mod tests { PkVectorSourceFile::new("f0".into(), 3).unwrap(), PkVectorSourceFile::new("f1".into(), 2).unwrap(), ]; - let live = build_live_row_ids(&files, &active_set(&["f0"]), &HashMap::new()) + let live = build_live_row_ids(&files, &active_set(&["f0"]), &HashMap::new(), None) .unwrap() .unwrap(); assert_eq!(live.iter().collect::>(), vec![0, 1, 2]); @@ -274,7 +326,7 @@ mod tests { let mut dvs = HashMap::new(); dvs.insert("f0".to_string(), dv(&[1])); // deletes global 1 dvs.insert("f1".to_string(), dv(&[0])); // deletes global 3 - let live = build_live_row_ids(&files, &active_set(&["f0", "f1"]), &dvs) + let live = build_live_row_ids(&files, &active_set(&["f0", "f1"]), &dvs, None) .unwrap() .unwrap(); assert_eq!(live.iter().collect::>(), vec![0, 2, 4]); @@ -290,6 +342,7 @@ mod tests { &meta, &active_set(&["f0", "f1"]), &HashMap::new(), + None, VectorSearchMetric::L2, ) .unwrap(); @@ -318,6 +371,7 @@ mod tests { &meta, &active_set(&["f0"]), &HashMap::new(), + None, VectorSearchMetric::L2, ) .unwrap_err(); @@ -333,6 +387,7 @@ mod tests { &meta, &active_set(&["f0"]), &HashMap::new(), + None, VectorSearchMetric::L2, ) .unwrap_err(); @@ -349,6 +404,7 @@ mod tests { &meta, &active_set(&["f0"]), &dvs, + None, VectorSearchMetric::L2, ) .unwrap_err(); @@ -399,6 +455,7 @@ mod tests { &active_set(&["f0", "f1"]), &dvs, &HashMap::new(), + None, ) .unwrap(); // Sorted BEST_FIRST by distance: (f1,0) dist 1.0 then (f0,0) dist 3.0. @@ -431,6 +488,7 @@ mod tests { &active_set(&["f0"]), &HashMap::new(), &HashMap::new(), + None, ) .unwrap_err(); assert!(err.to_string().contains("positive")); @@ -456,8 +514,148 @@ mod tests { &active_set(&["f0"]), &HashMap::new(), &HashMap::new(), + None, ) .unwrap(); assert!(results.is_empty()); } + + fn treemap(positions: &[u64]) -> roaring::RoaringTreemap { + let mut t = roaring::RoaringTreemap::new(); + for &p in positions { + t.insert(p); + } + t + } + + #[test] + fn test_build_live_row_ids_residual_intersects_with_active_and_dv() { + // f0 rows 0..3 (global 0,1,2), f1 rows 0..2 (global 3,4). Both active. + // dv on f0 deletes pos1 (global 1). residual allows f0={0,1}, f1 has no + // entry (empty allow). Result: f0 keeps {0} (1 is residual-allowed but + // deleted, 2 not residual-allowed); f1 contributes nothing. + let files = vec![ + PkVectorSourceFile::new("f0".into(), 3).unwrap(), + PkVectorSourceFile::new("f1".into(), 2).unwrap(), + ]; + let mut dvs = HashMap::new(); + dvs.insert("f0".to_string(), dv(&[1])); + let mut residual = HashMap::new(); + residual.insert("f0".to_string(), treemap(&[0, 1])); + let live = build_live_row_ids(&files, &active_set(&["f0", "f1"]), &dvs, Some(&residual)) + .unwrap() + .unwrap(); + assert_eq!(live.iter().collect::>(), vec![0]); + } + + #[test] + fn test_build_live_row_ids_residual_maps_positions_across_file_offsets() { + // f0 rows global 0,1,2; f1 rows global 3,4. residual allows f0={2}, f1={1}. + // f1 physical pos 1 -> global 3 + 1 = 4. Result {2, 4}. No DV. + let files = vec![ + PkVectorSourceFile::new("f0".into(), 3).unwrap(), + PkVectorSourceFile::new("f1".into(), 2).unwrap(), + ]; + let mut residual = HashMap::new(); + residual.insert("f0".to_string(), treemap(&[2])); + residual.insert("f1".to_string(), treemap(&[1])); + let live = build_live_row_ids( + &files, + &active_set(&["f0", "f1"]), + &HashMap::new(), + Some(&residual), + ) + .unwrap() + .unwrap(); + assert_eq!(live.iter().collect::>(), vec![2, 4]); + } + + #[test] + fn test_build_live_row_ids_residual_some_returns_mask_even_when_all_active_no_dv() { + // All active, no DV: without residual this returns None. With a residual + // present, a mask is always required. + let files = [PkVectorSourceFile::new("f0".into(), 3).unwrap()]; + let mut residual = HashMap::new(); + residual.insert("f0".to_string(), treemap(&[0, 2])); + let live = build_live_row_ids( + &files, + &active_set(&["f0"]), + &HashMap::new(), + Some(&residual), + ) + .unwrap() + .expect("residual present -> mask required"); + assert_eq!(live.iter().collect::>(), vec![0, 2]); + } + + #[test] + fn test_build_live_row_ids_rejects_out_of_range_residual_position() { + // Source file "f0" has 3 rows (valid positions 0..=2). A residual allow-list + // naming position 3 is out of range and must fail loud, not be skipped. + let files = source_meta(&[("f0", 3)]); + let mut residual = HashMap::new(); + residual.insert("f0".to_string(), treemap(&[0, 3])); + let err = build_live_row_ids( + files.source_files(), + &active_set(&["f0"]), + &HashMap::new(), + Some(&residual), + ) + .unwrap_err(); + assert!(err.to_string().contains("out of range")); + } + + #[test] + fn test_map_ann_results_rejects_hit_outside_residual_allow_list() { + // ordinal 1 -> (f0, 1). Residual allows only {0} in f0, so a hit at position 1 + // (e.g. an ANN reader that ignored include_row_ids) must fail loud. + let meta = source_meta(&[("f0", 3)]); + let mut residual = HashMap::new(); + residual.insert("f0".to_string(), treemap(&[0])); + let err = map_ann_results( + &[(1u64, 0.5)], + &meta, + &active_set(&["f0"]), + &HashMap::new(), + Some(&residual), + VectorSearchMetric::L2, + ) + .unwrap_err(); + assert!(err.to_string().contains("residual")); + } + + #[test] + fn test_vindex_adapter_sets_include_row_ids_to_residual_intersection() { + // Recording scorer captures the include_row_ids the adapter built. All + // active, no DV, residual f0={0,2} -> include_row_ids must equal {0,2}. + use std::sync::{Arc, Mutex}; + let seen_rows: Arc>>> = Arc::new(Mutex::new(None)); + let scorer_rows = Arc::clone(&seen_rows); + let searcher = VindexAnnSearcher::new( + "embedding".to_string(), + Box::new(move |_segment: &BucketAnnSegment, search: &VectorSearch| { + *scorer_rows.lock().unwrap() = search + .include_row_ids + .as_ref() + .map(|t| t.iter().collect::>()); + Ok(None) + }), + ); + let segment = BucketAnnSegment::for_test(source_meta(&[("f0", 3)])); + let mut residual = HashMap::new(); + residual.insert("f0".to_string(), treemap(&[0, 2])); + searcher + .search( + &segment, + &[0.0, 0.0], + VectorSearchMetric::L2, + 2, + &active_set(&["f0"]), + &HashMap::new(), + &HashMap::new(), + Some(&residual), + ) + .unwrap(); + assert_eq!(seen_rows.lock().unwrap().clone(), Some(vec![0, 2])); + } } diff --git a/crates/paimon/src/vindex/pkvector/bucket.rs b/crates/paimon/src/vindex/pkvector/bucket.rs index fcd7d1fed..604ac7dcd 100644 --- a/crates/paimon/src/vindex/pkvector/bucket.rs +++ b/crates/paimon/src/vindex/pkvector/bucket.rs @@ -142,6 +142,13 @@ pub(crate) fn covered_source_files( /// /// `ann_searcher` may be `None` only when there are no ANN segments; segments /// present with `None` is an error. +/// +/// `residual_ranges` (when `Some`) is a residual-predicate allow-list keyed by +/// data-file name whose value is the set of physical row positions in that file +/// that pass the predicate; only those rows may produce candidates. `None` applies +/// no residual restriction (every row is allowed). A file absent from the map (or +/// with an empty set) has no allowed rows and produces no candidates. Mirrors Java +/// `rowRangesByFile`. #[allow(clippy::too_many_arguments)] pub(crate) fn bucket_search( ann_searcher: Option<&dyn PkVectorAnnSearcher>, @@ -156,6 +163,7 @@ pub(crate) fn bucket_search( limit: usize, search_options: &HashMap, skip_exact_fallback: bool, + residual_ranges: Option<&HashMap>, ) -> crate::Result> { if limit == 0 { return Err(data_invalid("vector search limit must be positive")); @@ -240,6 +248,7 @@ pub(crate) fn bucket_search( &active_source_files, deletion_vectors, search_options, + residual_ranges, )? { add_candidate(&mut heap, result, limit); } @@ -250,13 +259,35 @@ pub(crate) fn bucket_search( if covered.contains(&file.file_name) { continue; } + // Residual allow-list: when present, only rows whose physical position + // passes the predicate may produce candidates. A file with no entry (or + // an empty entry) has no allowed rows, so it is skipped without reading. + let residual_allowed: Option<&roaring::RoaringTreemap> = match residual_ranges { + Some(ranges) => match ranges.get(&file.file_name) { + Some(allowed) if !allowed.is_empty() => Some(allowed), + _ => continue, + }, + None => None, + }; let dv = deletion_vectors.get(&file.file_name).cloned(); let is_excluded = move |position: i64| -> bool { - match &dv { + let dv_deleted = match &dv { Some(dv) => u64::try_from(position) .map(|p| dv.is_deleted(p)) .unwrap_or(false), None => false, + }; + if dv_deleted { + return true; + } + match residual_allowed { + // No residual restriction: the row is allowed. + None => false, + // Residual present: exclude positions outside the allow-list. + Some(allowed) => match u64::try_from(position) { + Ok(p) => !allowed.contains(p), + Err(_) => true, + }, } }; let mut reader = exact_reader_factory(file)?; @@ -319,6 +350,7 @@ mod tests { _active_source_files: &HashSet, _dvs: &HashMap>, _opts: &HashMap, + _residual_ranges: Option<&HashMap>, ) -> crate::Result> { Ok(self.result.clone()) } @@ -339,6 +371,7 @@ mod tests { 0, &HashMap::new(), false, + None, ) .unwrap_err(); assert!(err.to_string().contains("positive")); @@ -379,6 +412,7 @@ mod tests { 3, &HashMap::new(), false, + None, ) .unwrap(); // Top-3 BEST_FIRST: (data-1,0), (data-1,1), (data-1,2) — the larger @@ -394,10 +428,10 @@ mod tests { #[test] fn nan_ann_hit_never_evicts_finite_candidate_from_top1() { - // The core failure mode from review: an ANN hit with a negative-NaN - // distance must not win the single bucket Top-1 slot over a finite hit. - // Under f32::total_cmp the -NaN would rank best and evict the finite - // candidate here in the bucket heap, before any cross-bucket merge. + // The core failure mode: an ANN hit with a negative-NaN distance must not + // win the single bucket Top-1 slot over a finite hit. Under f32::total_cmp + // the -NaN would rank best and evict the finite candidate here in the + // bucket heap, before any cross-bucket merge. let negative_nan = f32::from_bits(0xffc00000); assert!(negative_nan.is_nan()); let segment = BucketAnnSegment::for_test(meta(&[("data-1", 2)])); @@ -428,6 +462,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap(); assert_eq!(results.len(), 1); @@ -467,6 +502,7 @@ mod tests { 2, &HashMap::new(), false, + None, ) .unwrap(); assert_eq!( @@ -516,6 +552,7 @@ mod tests { 2, &HashMap::new(), false, + None, ) .unwrap(); // Candidates: data-2 pos0 {1,0} dist 1.0; data-1 pos1 {2,0} dist 4.0. @@ -552,6 +589,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap_err(); assert!(err.to_string().contains("duplicate") || err.to_string().contains("Duplicate")); @@ -576,6 +614,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap_err(); assert!( @@ -614,6 +653,7 @@ mod tests { 2, &HashMap::new(), false, + None, ) .unwrap(); assert_eq!( @@ -644,6 +684,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap_err(); assert!( @@ -669,6 +710,7 @@ mod tests { 2, &HashMap::new(), true, // skip_exact_fallback + None, ) .unwrap(); assert!(results.is_empty()); @@ -702,6 +744,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap_err(); assert!( @@ -738,6 +781,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap_err(); assert!( @@ -763,6 +807,7 @@ mod tests { 1, &HashMap::new(), false, + None, ) .unwrap_err(); assert!(err.to_string().contains("row count") || err.to_string().contains("-1")); @@ -796,4 +841,162 @@ mod tests { let covered = covered_source_files(&[segment], &active); assert!(covered.is_empty()); } + + fn treemap(positions: &[u64]) -> roaring::RoaringTreemap { + let mut t = roaring::RoaringTreemap::new(); + for &p in positions { + t.insert(p); + } + t + } + + #[test] + fn test_exact_residual_allow_list_restricts_positions() { + // 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 = |_: &BucketActiveFile| -> crate::Result> { + 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]), + ], + ))) + }; + let mut residual: HashMap = HashMap::new(); + residual.insert("data-1".into(), treemap(&[0, 2])); + let results = bucket_search( + None, + &[], + &[active("data-1", 3)], + &HashMap::new(), + &mut factory, + &[0.0, 0.0], + VectorSearchMetric::L2, + 5, + &HashMap::new(), + false, + Some(&residual), + ) + .unwrap(); + assert_eq!( + results, + vec![ + PkVectorSearchResult { + data_file_name: "data-1".into(), + row_position: 0, + distance: 1.0 + }, + PkVectorSearchResult { + data_file_name: "data-1".into(), + row_position: 2, + distance: 9.0 + }, + ] + ); + } + + #[test] + fn test_exact_residual_file_absent_from_map_is_skipped_without_reading() { + // 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 = RefCell::new(Vec::::new()); + let mut factory = |f: &BucketActiveFile| -> crate::Result> { + calls.borrow_mut().push(f.file_name.clone()); + Ok(Box::new(ArrayReader::new( + 2, + vec![Some(vec![1.0, 0.0]), Some(vec![2.0, 0.0])], + ))) + }; + let mut residual: HashMap = HashMap::new(); + residual.insert("data-1".into(), treemap(&[0, 1])); + let results = bucket_search( + None, + &[], + &[active("data-1", 2), active("data-2", 2)], + &HashMap::new(), + &mut factory, + &[0.0, 0.0], + VectorSearchMetric::L2, + 5, + &HashMap::new(), + false, + Some(&residual), + ) + .unwrap(); + // Only data-1 rows appear; data-2 was never read. + assert!(results.iter().all(|r| r.data_file_name == "data-1")); + assert_eq!(calls.borrow().as_slice(), &["data-1".to_string()]); + } + + #[test] + fn test_exact_residual_empty_set_file_is_skipped_without_reading() { + // data-1 has an entry but it is empty -> no allowed rows, skipped without + // reading. Mirrors a file with no residual matches. + let calls = RefCell::new(0); + let mut factory = |_: &BucketActiveFile| -> crate::Result> { + *calls.borrow_mut() += 1; + unreachable!("data-1 has an empty allow set and must not be read") + }; + let mut residual: HashMap = HashMap::new(); + residual.insert("data-1".into(), treemap(&[])); + let results = bucket_search( + None, + &[], + &[active("data-1", 3)], + &HashMap::new(), + &mut factory, + &[0.0, 0.0], + VectorSearchMetric::L2, + 5, + &HashMap::new(), + false, + Some(&residual), + ) + .unwrap(); + assert!(results.is_empty()); + assert_eq!(*calls.borrow(), 0); + } + + #[test] + 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 = |_: &BucketActiveFile| -> crate::Result> { + 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]), + ], + ))) + }; + let mut dvs: HashMap> = HashMap::new(); + let mut bm = RoaringBitmap::new(); + bm.insert(0); // pos0 deleted + dvs.insert("data-1".into(), Arc::new(DeletionVector::from_bitmap(bm))); + let mut residual: HashMap = HashMap::new(); + residual.insert("data-1".into(), treemap(&[0, 1, 2])); + let results = bucket_search( + None, + &[], + &[active("data-1", 3)], + &dvs, + &mut factory, + &[0.0, 0.0], + VectorSearchMetric::L2, + 5, + &HashMap::new(), + false, + Some(&residual), + ) + .unwrap(); + assert_eq!( + results.iter().map(|r| r.row_position).collect::>(), + vec![1, 2] + ); + } } diff --git a/crates/paimon/tests/pk_vector_baseline_test.rs b/crates/paimon/tests/pk_vector_baseline_test.rs index e1e3177a7..55b144a83 100644 --- a/crates/paimon/tests/pk_vector_baseline_test.rs +++ b/crates/paimon/tests/pk_vector_baseline_test.rs @@ -58,6 +58,9 @@ use paimon::spec::{ DataFileMeta, DataType, FloatType, GlobalIndexMeta, IndexFileMeta, IntType, Schema, TableSchema, VectorType, }; +// Used only by the residual end-to-end test, which is gated off Windows. +#[cfg(not(windows))] +use paimon::spec::{Datum, Predicate, PredicateBuilder}; use paimon::table::{CommitMessage, SchemaManager, Table, TableCommit}; use paimon_vindex_core::index::{VectorIndexConfig, VectorIndexTrainer, VectorIndexWriter}; use paimon_vindex_core::io::PosWriter; @@ -103,9 +106,19 @@ fn analytic_topk(query: &[f32], vectors: &[[f32; DIM]], k: usize) -> Vec<(u64, f /// Table options that route searches into the primary-key vector branch /// (`VectorSearchBuilder::execute_primary_key_vector_search`). Default search /// mode is FAST, so only the ANN segment is consulted (no exact fallback). +/// +/// `deletion-vectors.enabled = true` (and merge-on-read left at its default +/// `false`) is what makes the table expose physical rows directly, the +/// precondition Java `PrimaryKeyVectorScan` requires before a residual data +/// predicate (`with_filter`) may be applied post-recall. This fixture never +/// actually deletes a row, so no deletion files are written; the option only +/// satisfies the residual guard and does not otherwise change the write/read +/// path (default `deduplicate` merge-engine, no deletion files -> no DV factory +/// is built on read). fn table_options() -> Vec<(String, String)> { vec![ ("bucket".to_string(), "1".to_string()), + ("deletion-vectors.enabled".to_string(), "true".to_string()), ( "pk-vector.index.columns".to_string(), VECTOR_COLUMN.to_string(), @@ -685,3 +698,229 @@ async fn pk_vector_read_orders_rows_best_first_not_by_position() { ); } } + +/// Fixture #3 (residual): the unrestricted nearest neighbours sit at low ids +/// (0, 1, 2), but the residual predicate `id >= 3` excludes exactly those, so +/// the residual result set is disjoint from the unfiltered one. Among the rows +/// the residual keeps (ids 3, 4, 5) the best-first order is [4, 5, 3], which is +/// neither the ascending id order [3, 4, 5] nor a prefix of the unfiltered +/// order — proving the residual genuinely reshapes the result. +/// +/// query [10,0,0,0] +/// pos0 [10,0,0,0] -> 0 (unfiltered nearest; excluded by id >= 3) +/// pos1 [ 9,0,0,0] -> 1 (excluded) +/// pos2 [ 8,0,0,0] -> 4 (excluded) +/// pos3 [ 5,0,0,0] -> 25 (kept; farthest of the kept rows) +/// pos4 [ 7,0,0,0] -> 9 (kept; nearest of the kept rows) +/// pos5 [ 6,0,0,0] -> 16 (kept) +/// Strict gaps 0 < 1 < 4 < 9 < 16 < 25 make every top-k order unique. +/// unfiltered top-3 = [0, 1, 2] +/// residual (id >= 3) top-3 = [4, 5, 3] +#[cfg(not(windows))] +fn fixture_residual() -> ([f32; DIM], Vec<[f32; DIM]>) { + let query = [10.0, 0.0, 0.0, 0.0]; + let vectors = vec![ + [10.0, 0.0, 0.0, 0.0], // pos 0 -> 0 + [9.0, 0.0, 0.0, 0.0], // pos 1 -> 1 + [8.0, 0.0, 0.0, 0.0], // pos 2 -> 4 + [5.0, 0.0, 0.0, 0.0], // pos 3 -> 25 + [7.0, 0.0, 0.0, 0.0], // pos 4 -> 9 + [6.0, 0.0, 0.0, 0.0], // pos 5 -> 16 + ]; + (query, vectors) +} + +/// Run `execute_read()` with a residual `filter` attached via `with_filter` and +/// flatten the stream into per-row `(id, score)` tuples in emission order +/// (best-first), returning the collected batches too for schema / row-content +/// assertions. Mirrors `read_id_and_scores` but exercises the residual path. +#[cfg(not(windows))] +async fn read_id_and_scores_filtered( + table: &Table, + query: Vec, + limit: usize, + filter: Predicate, +) -> (Vec, Vec, Vec) { + let mut builder = table.new_vector_search_builder(); + builder + .with_vector_column(VECTOR_COLUMN) + .with_query_vector(query) + .with_limit(limit) + .with_filter(filter); + let batches = builder + .execute_read() + .await + .expect("primary-key vector residual read failed") + .try_collect::>() + .await + .expect("collecting residual read batches failed"); + + let ids: Vec = batches + .iter() + .flat_map(|b| { + let idx = b.schema().index_of("id").unwrap(); + b.column(idx) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }) + .collect(); + let scores: Vec = batches + .iter() + .flat_map(|b| { + let idx = b.schema().index_of("_PKEY_VECTOR_SCORE").unwrap(); + b.column(idx) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }) + .collect(); + (ids, scores, batches) +} + +/// End-to-end coverage of the residual data predicate (`with_filter`) on the +/// public primary-key vector path. Mirrors Java `PrimaryKeyVectorRead`'s +/// residual-filter support: the residual columns are re-read per candidate file, +/// the surviving physical positions are folded into recall, and only rows +/// satisfying the predicate remain — best-first and Top-K preserved. +/// +/// The fixture is built so the residual actually changes the result set: the +/// unfiltered top-3 is [0, 1, 2] while the residual (`id >= 3`) top-3 is +/// [4, 5, 3]. The two sets are disjoint, so a residual that silently did nothing +/// (or was ignored) would surface here. +// Gated off Windows for the same `file://` tempdir reason as the tests above. +#[cfg(not(windows))] +#[tokio::test] +async fn pk_vector_residual_filter_excludes_non_matching_rows() { + let (query, vectors) = fixture_residual(); + let (_tmp, table) = build_table(&query, &vectors, 3).await; + + // Ground truth: unrestricted top-3, and the top-3 restricted to ids >= 3 + // (the residual only ranks rows whose id passes the predicate). + let unfiltered = analytic_topk(&query, &vectors, 3); + let unfiltered_ids: Vec = unfiltered.iter().map(|(id, _)| *id).collect(); + assert_eq!( + unfiltered_ids, + vec![0, 1, 2], + "fixture guard: unfiltered top-3 must be [0, 1, 2]" + ); + + let residual_threshold = 3; + let mut residual_ranked: Vec<(u64, f32)> = analytic_topk(&query, &vectors, vectors.len()) + .into_iter() + .filter(|(id, _)| *id >= residual_threshold) + .collect(); + residual_ranked.truncate(3); + let expected_ids: Vec = residual_ranked.iter().map(|(id, _)| *id as i32).collect(); + let expected_scores: Vec = residual_ranked.iter().map(|(_, d)| l2_score(*d)).collect(); + // The whole point of this fixture: residual result != unfiltered result, and + // best-first over the kept rows is not ascending id order. + assert_eq!( + expected_ids, + vec![4, 5, 3], + "fixture guard: residual (id >= 3) top-3 must be best-first [4, 5, 3]" + ); + + // Build the residual predicate on the data column `id` via the public + // PredicateBuilder, exactly as a caller would. + let residual = PredicateBuilder::new(table.schema().fields()) + .greater_or_equal("id", Datum::Int(residual_threshold as i32)) + .expect("build residual predicate on id"); + + // Search-only: unfiltered vs residual must differ, and every residual hit + // must satisfy the predicate (id >= 3), disjoint from the unfiltered set. + let unfiltered_result = table + .new_vector_search_builder() + .with_vector_column(VECTOR_COLUMN) + .with_query_vector(query.to_vec()) + .with_limit(3) + .execute_scored() + .await + .expect("unfiltered primary-key vector search failed"); + assert_eq!( + unfiltered_result.row_ids, unfiltered_ids, + "unfiltered search must return [0, 1, 2]" + ); + + let residual_result = table + .new_vector_search_builder() + .with_vector_column(VECTOR_COLUMN) + .with_query_vector(query.to_vec()) + .with_limit(3) + .with_filter(residual.clone()) + .execute_scored() + .await + .expect("residual primary-key vector search failed"); + let residual_row_ids: Vec = expected_ids.iter().map(|&id| id as u64).collect(); + assert_eq!( + residual_result.row_ids, residual_row_ids, + "residual search must return best-first [4, 5, 3]" + ); + assert_ne!( + residual_result.row_ids, unfiltered_result.row_ids, + "residual must change the result set relative to no filter" + ); + for &id in &residual_result.row_ids { + assert!( + id >= residual_threshold, + "residual search returned id {id} that fails the predicate id >= {residual_threshold}" + ); + } + + // Search-and-read with the residual: default projection materializes id + + // vector column, best-first, with an aligned `_PKEY_VECTOR_SCORE`. + let (ids, scores, batches) = + read_id_and_scores_filtered(&table, query.to_vec(), 3, residual).await; + + assert_eq!( + ids, expected_ids, + "residual read must emit only kept rows, best-first [4, 5, 3]" + ); + for &id in &ids { + assert!( + id >= residual_threshold as i32, + "residual read returned id {id} that fails the predicate id >= {residual_threshold}" + ); + } + + // Row content: the materialized vector for each emitted row equals the source + // vector at that physical position. + let got_vectors = collect_vectors(&batches); + assert_eq!(got_vectors.len(), 3, "three kept rows expected"); + for (row_idx, (id, _)) in residual_ranked.iter().enumerate() { + assert_eq!( + got_vectors[row_idx], + vectors[*id as usize].to_vec(), + "materialized vector for row id {id} diverges from source data" + ); + } + + // Score alignment on the residual read path. + assert_eq!(scores.len(), 3); + for (got, want) in scores.iter().zip(&expected_scores) { + assert!( + (got - want).abs() < 1e-4, + "residual read score diverges: got {got}, want {want}" + ); + } + + // Hidden metadata columns must not leak, even on the residual path. + for batch in &batches { + assert!( + batch.schema().index_of("_ROW_ID").is_err(), + "_ROW_ID must not leak into residual read output" + ); + assert!( + batch.schema().index_of("_PKEY_VECTOR_POSITION").is_err(), + "_PKEY_VECTOR_POSITION must not leak into residual read output" + ); + assert!( + batch.schema().index_of(VECTOR_COLUMN).is_ok(), + "default projection must materialize the vector column" + ); + } +}