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 595e3580 feat(table): add primary-key hybrid (vector + full-text)
search (#568)
595e3580 is described below
commit 595e3580e1f4e15d302d60f6726278d0245ae788
Author: Junrui Lee <[email protected]>
AuthorDate: Sun Jul 26 19:42:05 2026 +0800
feat(table): add primary-key hybrid (vector + full-text) search (#568)
---
crates/paimon/src/table/hybrid_search_builder.rs | 1355 +++++++++++++++++++++-
crates/paimon/src/table/mod.rs | 2 +
crates/paimon/src/table/pk_full_text_read.rs | 138 ++-
crates/paimon/src/table/pk_full_text_scan.rs | 29 +-
crates/paimon/src/table/pk_search_position.rs | 250 ++++
crates/paimon/src/table/pk_search_ranker.rs | 452 ++++++++
crates/paimon/src/table/pk_vector_scan.rs | 40 +-
crates/paimon/src/table/vector_search_builder.rs | 323 +++++-
8 files changed, 2548 insertions(+), 41 deletions(-)
diff --git a/crates/paimon/src/table/hybrid_search_builder.rs
b/crates/paimon/src/table/hybrid_search_builder.rs
index 11b7a026..ab290820 100644
--- a/crates/paimon/src/table/hybrid_search_builder.rs
+++ b/crates/paimon/src/table/hybrid_search_builder.rs
@@ -19,12 +19,33 @@
//!
//! Reference: `org.apache.paimon.table.source.HybridSearchBuilder`.
-use std::collections::HashMap;
+use std::collections::{BTreeMap, HashMap};
-use crate::spec::CoreOptions;
-use crate::table::{RowRange, Table};
+use arrow_array::RecordBatch;
+use futures::{stream, TryStreamExt};
+
+use crate::spec::{CoreOptions, SCAN_SNAPSHOT_ID_OPTION};
+use crate::table::data_file_reader::DataFileReader;
+use crate::table::pk_search_position::PrimaryKeySearchPosition;
+use crate::table::pk_search_ranker::{self, Ranking};
+use crate::table::pk_vector_indexed_split_read::{PkVectorIndexedSplit,
PkVectorIndexedSplitRead};
+use crate::table::pk_vector_orchestrator::build_indexed_splits;
+use crate::table::source::DataSplit;
+use crate::table::vector_search_builder::{
+ collect_ranked_rows, ensure_no_reserved_read_columns,
reorder_and_strip_position, RankedRow,
+};
+use crate::table::{ArrowRecordBatchStream, RowRange, Table};
use crate::vector_search::SearchResult;
+#[cfg(feature = "fulltext")]
+use crate::spec::GlobalIndexSearchMode;
+#[cfg(feature = "fulltext")]
+use crate::table::find_field_id_by_name;
+#[cfg(feature = "fulltext")]
+use crate::table::pk_full_text_read::{build_full_text_indexed_splits,
PrimaryKeyFullTextRead};
+#[cfg(feature = "fulltext")]
+use crate::table::pk_full_text_scan::PrimaryKeyFullTextScan;
+
const RRF_K: f32 = 60.0;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -265,7 +286,8 @@ impl<'a> HybridSearchBuilder<'a> {
}
pub async fn execute_scored(&self) -> crate::Result<SearchResult> {
-
CoreOptions::new(self.table.schema().options()).ensure_read_authorized()?;
+ let core = CoreOptions::new(self.table.schema().options());
+ core.ensure_read_authorized()?;
let limit = self.limit.ok_or_else(|| crate::Error::ConfigInvalid {
message: "Limit must be set via with_limit()".to_string(),
})?;
@@ -275,6 +297,23 @@ impl<'a> HybridSearchBuilder<'a> {
});
}
+ // A primary-key hybrid fuses PHYSICAL positions, not global row ids,
so a
+ // scored/row-range result is unsupported on it: fail loud and direct
callers
+ // to the materialized `execute_read`. A mixed PK/global set of routes
cannot
+ // be fused at all (different address spaces) and also fails loud. The
+ // append/data-evolution global path is unchanged (mirrors Java
`rank`).
+ match self.classify_routes(&core)? {
+ HybridAddressSpace::PrimaryKey => {
+ return Err(crate::Error::DataInvalid {
+ message: "primary-key hybrid search does not produce
global row ids; use the \
+ materialized read (execute_read) instead"
+ .to_string(),
+ source: None,
+ });
+ }
+ HybridAddressSpace::Global => {}
+ }
+
let mut route_results = Vec::with_capacity(self.routes.len());
for route in &self.routes {
let result = match route.kind {
@@ -301,6 +340,560 @@ impl<'a> HybridSearchBuilder<'a> {
Ok(rank_results(self.ranker, &route_results, limit))
}
+
+ /// Materialize the fused hybrid search hits into Arrow rows,
best-fused-score
+ /// first, with a `__paimon_search_score` column appended. Only the
primary-key
+ /// hybrid path can materialize rows: every route must resolve to a
physical
+ /// primary-key route (vector column configured in
`pk-vector.index.columns`,
+ /// full-text column in `pk-full-text.index.columns` with data-evolution
+ /// disabled). A mixed PK/global set of routes fails loud, and an
all-global
+ /// (append/data-evolution) hybrid is unsupported here — those use
+ /// `execute`/`execute_scored`. Mirrors Java `HybridSearchBuilderImpl` PK
path.
+ pub async fn execute_read(&self) -> crate::Result<ArrowRecordBatchStream> {
+ let core = CoreOptions::new(self.table.schema().options());
+ core.ensure_read_authorized()?;
+ let limit = self.limit.ok_or_else(|| crate::Error::ConfigInvalid {
+ message: "Limit must be set via with_limit()".to_string(),
+ })?;
+ if self.routes.is_empty() {
+ return Err(crate::Error::ConfigInvalid {
+ message: "Routes cannot be empty".to_string(),
+ });
+ }
+
+ match self.classify_routes(&core)? {
+ HybridAddressSpace::PrimaryKey => {
+ self.execute_primary_key_hybrid_read(&core, limit).await
+ }
+ HybridAddressSpace::Global => Err(crate::Error::Unsupported {
+ message: "materialized hybrid read (execute_read) is only
supported on the \
+ primary-key hybrid path; use execute/execute_scored
for the \
+ append/data-evolution path"
+ .to_string(),
+ }),
+ }
+ }
+
+ /// Classify the route set as an all-primary-key or an all-global hybrid,
+ /// failing loud when the two address spaces are mixed (mirrors Java
`rank`).
+ fn classify_routes(&self, core: &CoreOptions<'_>) ->
crate::Result<HybridAddressSpace> {
+ let mut any_pk = false;
+ let mut any_global = false;
+ for route in &self.routes {
+ if route_is_primary_key(core, route) {
+ any_pk = true;
+ } else {
+ any_global = true;
+ }
+ }
+ if any_pk && any_global {
+ return Err(crate::Error::Unsupported {
+ message: "Hybrid search cannot combine physical primary-key
positions and global \
+ row-id address spaces."
+ .to_string(),
+ });
+ }
+ Ok(if any_pk {
+ HybridAddressSpace::PrimaryKey
+ } else {
+ HybridAddressSpace::Global
+ })
+ }
+
+ /// The all-primary-key hybrid read: run each route's candidate producer
against
+ /// the ONE pinned snapshot, convert its candidates to shared physical
positions,
+ /// fuse them via the configured ranker (each route uses its own limit, the
+ /// fusion uses the builder limit), select the physical source files the
fused
+ /// positions reference, build indexed splits carrying the FUSED raw
scores,
+ /// materialize, and reorder best-fused-first while stripping the internal
+ /// position column.
+ async fn execute_primary_key_hybrid_read(
+ &self,
+ core: &CoreOptions<'_>,
+ limit: usize,
+ ) -> crate::Result<ArrowRecordBatchStream> {
+ // The materialized read projects every user column and then appends
the
+ // internal `_PKEY_VECTOR_POSITION` and `__paimon_search_score`
columns; a
+ // user column whose name collides with either (or with `_ROW_ID`)
would be
+ // shadowed by a positional lookup, corrupting the reorder/strip and
the
+ // fused score. Reject it up front — before any route runs and even
when the
+ // fused result is empty — reusing the primary-key vector read's guard.
+ ensure_no_reserved_read_columns(self.table.schema().fields())?;
+
+ // Snapshot pinning: resolve ONE snapshot for the whole primary-key
hybrid
+ // read and plan every route against it, mirroring Java
+ // `HybridSearchBuilderImpl.routeBuilders()` (which resolves the
snapshot
+ // once and injects it into every route builder). Only the "read
latest"
+ // case is racy — a concurrent commit landing between the two route
plans
+ // would otherwise pick different snapshots — because every time-travel
+ // selector resolves deterministically to the same snapshot on each
plan.
+ // So pin latest once via `scan.snapshot-id` and leave the
already-pinned
+ // (time-travel) paths to resolve their own fixed snapshot.
+ let pinned_table = self.resolve_pinned_route_table().await?;
+ let route_table: &Table = pinned_table.as_ref().unwrap_or(self.table);
+
+ // Per-route search: positions (converted from candidates) +
single-file
+ // source splits + the snapshot each route's plan pinned.
+ let mut routes: Vec<PkRoute> = Vec::with_capacity(self.routes.len());
+ for route in &self.routes {
+ let pk_route = match route.kind {
+ HybridSearchRouteKind::Vector => {
+ self.pk_vector_route(route_table, core, route).await?
+ }
+ HybridSearchRouteKind::FullText => {
+ self.pk_full_text_route(route_table, core, route).await?
+ }
+ };
+ routes.push(pk_route);
+ }
+
+ // Snapshot pinning: every route's plan must resolve the same
snapshot, so
+ // never fuse across versions (mirrors Java `rankPhysical`'s mismatch
guard).
+ // With a pinned snapshot above this is a defensive guard, not the
primary
+ // mechanism.
+ check_single_snapshot(routes.iter().map(|r| r.snapshot_id))?;
+
+ // Fuse the per-route rankings into a single best-first position list
bounded
+ // by the builder limit.
+ let fused = self.fuse_positions(&routes, limit)?;
+ if fused.is_empty() {
+ return Ok(Box::pin(stream::empty()));
+ }
+
+ // Physical sources: collect the single-file source splits across
routes,
+ // validating cross-route metadata/DV consistency, then select only
the files
+ // the fused positions reference (mirrors Java `physicalSources`).
+ let available = collect_physical_sources(&routes)?;
+
+ // Rank each fused position by its best-first ordinal via its FULL
physical
+ // key so the file/position materialization order can be reduced back
to
+ // best-first.
+ let mut rank_of: HashMap<(Vec<u8>, i32, String, i64), usize> =
HashMap::new();
+ for (rank, position) in fused.iter().enumerate() {
+ rank_of.insert(
+ (
+ position.partition().to_serialized_bytes(),
+ position.bucket(),
+ position.data_file_name().to_string(),
+ position.row_position(),
+ ),
+ rank,
+ );
+ }
+
+ // Build the materialization splits DIRECTLY from the fused positions,
+ // carrying the FUSED raw scores (no distance conversion): those are
the
+ // final relevance scores the output must expose.
+ let indexed_splits = build_hybrid_indexed_splits(&fused, &available)?;
+
+ // A predicate-free materialization reader projecting every user
column; the
+ // indexed-split read appends the score column itself.
+ let materialize_reader = DataFileReader::new(
+ self.table.file_io().clone(),
+ self.table.schema_manager().clone(),
+ self.table.schema().id(),
+ self.table.schema().fields().to_vec(),
+ self.table.schema().fields().to_vec(),
+ Vec::new(),
+ );
+
+ let mut batches: Vec<RecordBatch> = Vec::new();
+ let mut ranked: Vec<RankedRow> = Vec::new();
+ for indexed in indexed_splits {
+ let partition_bytes =
indexed.split.partition().to_serialized_bytes();
+ let bucket = indexed.split.bucket();
+ let file_name = indexed.split.data_files()[0].file_name.clone();
+ let mut read_stream =
+
PkVectorIndexedSplitRead::new(materialize_reader.clone()).read(&indexed)?;
+ while let Some(batch) = read_stream.try_next().await? {
+ let batch_index = batches.len();
+ collect_ranked_rows(
+ &batch,
+ batch_index,
+ &partition_bytes,
+ bucket,
+ &file_name,
+ &rank_of,
+ &mut ranked,
+ )?;
+ batches.push(batch);
+ }
+ }
+
+ let output = reorder_and_strip_position(&batches, ranked)?;
+ Ok(Box::pin(stream::iter(output.into_iter().map(Ok))))
+ }
+
+ /// Resolve the ONE snapshot every primary-key route must plan against, as
an
+ /// optional pinned table copy. Mirrors Java
+ /// `HybridSearchBuilderImpl.routeBuilders()`, which resolves a single
snapshot
+ /// up front and injects it into every route builder.
+ ///
+ /// A time-travel selector (`scan.version` / `scan.timestamp-millis` /
+ /// `scan.snapshot-id` / `scan.tag-name`) already resolves
deterministically to
+ /// the same snapshot on every route plan, so no extra pinning is needed —
+ /// return `None` and let each route resolve it. Only the default "read
latest"
+ /// path is racy: a concurrent commit landing between two route plans
would let
+ /// them resolve different snapshots. For that path resolve the latest
snapshot
+ /// id ONCE and return a table copy pinned to it via `scan.snapshot-id`,
so both
+ /// routes plan the same version. A table with no snapshot at all also
returns
+ /// `None` (nothing to pin; every route plans empty).
+ async fn resolve_pinned_route_table(&self) -> crate::Result<Option<Table>>
{
+ let core = CoreOptions::new(self.table.schema().options());
+ // Already targeting a fixed snapshot (resolved travel copy or a
selector
+ // that resolves deterministically): every route agrees without
pinning.
+ if self.table.has_resolved_travel_snapshot() ||
core.try_time_travel_selector()?.is_some() {
+ return Ok(None);
+ }
+ // Read-latest: pin the current latest snapshot once so a concurrent
commit
+ // cannot split the routes across versions.
+ let Some(latest) = self
+ .table
+ .snapshot_manager()
+ .get_latest_snapshot_id()
+ .await?
+ else {
+ return Ok(None);
+ };
+ let pinned = self.table.copy_with_options(HashMap::from([(
+ SCAN_SNAPSHOT_ID_OPTION.to_string(),
+ latest.to_string(),
+ )]));
+ Ok(Some(pinned))
+ }
+
+ /// Run the vector route's primary-key candidate producer and convert its
hits
+ /// into shared physical positions (distance → score via the resolved
metric),
+ /// keeping the route's single-file source splits and pinned snapshot.
+ async fn pk_vector_route(
+ &self,
+ table: &Table,
+ core: &CoreOptions<'_>,
+ route: &HybridSearchRoute,
+ ) -> crate::Result<PkRoute> {
+ let vector = route.vector.as_deref().expect("validated vector route");
+ let mut builder = table.new_vector_search_builder();
+ builder
+ .with_vector_column(&route.field_name)
+ .with_query_vector(vector.to_vec())
+ .with_limit(route.limit)
+ .with_options(route.options.clone());
+ let result = builder
+ .search_pk_route(core, &route.field_name, vector, route.limit)
+ .await?;
+ let positions = result
+ .candidates
+ .iter()
+ .map(|candidate| {
+ PrimaryKeySearchPosition::from_vector_candidate(candidate,
result.metric)
+ })
+ .collect::<crate::Result<Vec<_>>>()?;
+ let source_splits = build_indexed_splits(result.candidates,
&result.splits, result.metric)?
+ .into_iter()
+ .map(|split| split.split)
+ .collect();
+ Ok(PkRoute {
+ positions,
+ source_splits,
+ snapshot_id: result.snapshot_id,
+ weight: route.weight as f64,
+ })
+ }
+
+ #[cfg(feature = "fulltext")]
+ async fn pk_full_text_route(
+ &self,
+ table: &Table,
+ core: &CoreOptions<'_>,
+ route: &HybridSearchRoute,
+ ) -> crate::Result<PkRoute> {
+ // FAST-only: the primary-key full-text read searches
compaction-visible
+ // payloads and rejects FULL/DETAIL loud rather than silently
degrading,
+ // mirroring `full_text_search_builder::execute_read` and Java
+ // `PrimaryKeyFullTextRead.checkFastSearchMode`.
+ if core.global_index_search_mode()? != GlobalIndexSearchMode::Fast {
+ return Err(crate::Error::DataInvalid {
+ message: "primary-key full-text search supports only the FAST
global-index search \
+ mode"
+ .to_string(),
+ source: None,
+ });
+ }
+
+ let query = route
+ .full_text_query
+ .as_deref()
+ .expect("validated full-text route");
+ let field_id = find_field_id_by_name(table.schema().fields(),
&route.field_name)
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: format!(
+ "full-text search column '{}' does not exist",
+ route.field_name
+ ),
+ source: None,
+ })?;
+ let plan = PrimaryKeyFullTextScan::new(table, field_id, None)
+ .plan()
+ .await?;
+ let materialize_reader = DataFileReader::new(
+ table.file_io().clone(),
+ table.schema_manager().clone(),
+ table.schema().id(),
+ table.schema().fields().to_vec(),
+ table.schema().fields().to_vec(),
+ Vec::new(),
+ );
+ let read = PrimaryKeyFullTextRead::new(
+ table.file_io().clone(),
+ materialize_reader,
+ table.location().trim_end_matches('/').to_string(),
+ );
+ let result = read.search_route(&plan, query, route.limit).await?;
+ let positions = result
+ .candidates
+ .iter()
+ .map(PrimaryKeySearchPosition::from_full_text_candidate)
+ .collect::<crate::Result<Vec<_>>>()?;
+ let source_splits = build_full_text_indexed_splits(result.candidates,
result.splits)?
+ .into_iter()
+ .map(|split| split.split)
+ .collect();
+ Ok(PkRoute {
+ positions,
+ source_splits,
+ snapshot_id: result.snapshot_id,
+ weight: route.weight as f64,
+ })
+ }
+
+ #[cfg(not(feature = "fulltext"))]
+ async fn pk_full_text_route(
+ &self,
+ _table: &Table,
+ _core: &CoreOptions<'_>,
+ _route: &HybridSearchRoute,
+ ) -> crate::Result<PkRoute> {
+ Err(crate::Error::ConfigInvalid {
+ message: "primary-key full-text hybrid routes require the fulltext
feature".to_string(),
+ })
+ }
+
+ /// Fuse the per-route physical rankings via the configured ranker. Each
route
+ /// contributes a weighted ranking only when it has positions (mirrors Java
+ /// `rankPhysical`), and the fusion is bounded to the builder limit.
+ fn fuse_positions(
+ &self,
+ routes: &[PkRoute],
+ limit: usize,
+ ) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ let mut rankings = Vec::with_capacity(routes.len());
+ for route in routes {
+ if !route.positions.is_empty() {
+ rankings.push(Ranking::new(route.positions.clone(),
route.weight)?);
+ }
+ }
+ match self.ranker {
+ HybridSearchRanker::Rrf =>
pk_search_ranker::weighted_rrf(&rankings, limit),
+ HybridSearchRanker::WeightedScore =>
pk_search_ranker::weighted_score(&rankings, limit),
+ HybridSearchRanker::Mrr =>
pk_search_ranker::weighted_mrr(&rankings, limit),
+ }
+ }
+}
+
+/// One primary-key route's fusion input: its converted physical positions, the
+/// single-file source splits its hits reference (for physical-sources
+/// materialization), the snapshot its plan pinned, and its fusion weight.
+struct PkRoute {
+ positions: Vec<PrimaryKeySearchPosition>,
+ source_splits: Vec<DataSplit>,
+ snapshot_id: i64,
+ weight: f64,
+}
+
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+enum HybridAddressSpace {
+ PrimaryKey,
+ Global,
+}
+
+/// Whether a route resolves to a physical primary-key route. Vector: the
queried
+/// column is a configured `pk-vector.index.columns` entry (membership via the
+/// non-erroring accessor so a malformed config cannot misclassify an unrelated
+/// route). Full-text: data-evolution disabled AND the column configured in
+/// `pk-full-text.index.columns` (only when the fulltext feature is built).
+fn route_is_primary_key(core: &CoreOptions<'_>, route: &HybridSearchRoute) ->
bool {
+ match route.kind {
+ HybridSearchRouteKind::Vector => {
+ core.primary_key_vector_index_enabled()
+ && core
+ .primary_key_vector_index_columns()
+ .ok()
+ .is_some_and(|cols| cols.iter().any(|c| c ==
route.field_name()))
+ }
+ HybridSearchRouteKind::FullText => {
+ #[cfg(feature = "fulltext")]
+ {
+ !core.data_evolution_enabled()
+ && core
+ .primary_key_full_text_index_columns()
+ .iter()
+ .any(|c| c == route.field_name())
+ }
+ #[cfg(not(feature = "fulltext"))]
+ {
+ let _ = core;
+ false
+ }
+ }
+ }
+}
+
+/// Assert every primary-key route pinned the same snapshot, returning it.
Mirrors
+/// the `checkArgument` in Java `rankPhysical`: fusing across snapshots would
+/// search/materialize physical rows against inconsistent versions.
+fn check_single_snapshot(
+ snapshot_ids: impl IntoIterator<Item = i64>,
+) -> crate::Result<Option<i64>> {
+ let mut pinned: Option<i64> = None;
+ for snapshot_id in snapshot_ids {
+ match pinned {
+ None => pinned = Some(snapshot_id),
+ Some(existing) if existing != snapshot_id => {
+ return Err(crate::Error::DataInvalid {
+ message: format!(
+ "primary-key hybrid routes must use the same snapshot,
but found {existing} and {snapshot_id}"
+ ),
+ source: None,
+ });
+ }
+ _ => {}
+ }
+ }
+ Ok(pinned)
+}
+
+/// Physical file identity for the physical-sources map: `(partition bytes,
bucket,
+/// data file name)`.
+type PhysicalFileKey = (Vec<u8>, i32, String);
+
+/// Collect the single-file source splits across all routes into a map keyed by
+/// physical file identity, validating that a file appearing in more than one
route
+/// carries consistent metadata and deletion-file state. Mirrors Java
+/// `physicalSources`'s duplicate-consistency guard.
+fn collect_physical_sources(
+ routes: &[PkRoute],
+) -> crate::Result<HashMap<PhysicalFileKey, DataSplit>> {
+ let mut available: HashMap<PhysicalFileKey, DataSplit> = HashMap::new();
+ for route in routes {
+ for split in &route.source_splits {
+ if split.data_files().len() != 1 {
+ return Err(crate::Error::DataInvalid {
+ message: "primary-key scored source split must contain
exactly one data file"
+ .to_string(),
+ source: None,
+ });
+ }
+ let file = &split.data_files()[0];
+ let key = (
+ split.partition().to_serialized_bytes(),
+ split.bucket(),
+ file.file_name.clone(),
+ );
+ match available.get(&key) {
+ None => {
+ available.insert(key, split.clone());
+ }
+ Some(previous) => {
+ let prev_file = &previous.data_files()[0];
+ let prev_dv = previous
+ .data_deletion_files()
+ .and_then(|dfs| dfs.first().cloned().flatten());
+ let cur_dv = split
+ .data_deletion_files()
+ .and_then(|dfs| dfs.first().cloned().flatten());
+ let consistent = previous.snapshot_id() ==
split.snapshot_id()
+ && previous.bucket_path() == split.bucket_path()
+ && previous.total_buckets() == split.total_buckets()
+ && prev_file.file_size == file.file_size
+ && prev_file.row_count == file.row_count
+ && prev_dv == cur_dv;
+ if !consistent {
+ return Err(crate::Error::DataInvalid {
+ message: format!(
+ "primary-key hybrid routes contain
inconsistent metadata for data file {}",
+ file.file_name
+ ),
+ source: None,
+ });
+ }
+ }
+ }
+ }
+ }
+ Ok(available)
+}
+
+/// Group the fused positions by physical file, look up each file's single-file
+/// source split, and build one `PkVectorIndexedSplit` per file carrying the
FUSED
+/// raw scores aligned to ascending position order. A fused position
referencing a
+/// file no route sourced fails loud (mirrors Java `physicalSources`'s
selection).
+fn build_hybrid_indexed_splits(
+ fused: &[PrimaryKeySearchPosition],
+ available: &HashMap<PhysicalFileKey, DataSplit>,
+) -> crate::Result<Vec<PkVectorIndexedSplit>> {
+ // BTreeMap keeps a deterministic ascending group order. Value: (position,
score).
+ let mut groups: BTreeMap<PhysicalFileKey, Vec<(i64, f32)>> =
BTreeMap::new();
+ for position in fused {
+ let key = (
+ position.partition().to_serialized_bytes(),
+ position.bucket(),
+ position.data_file_name().to_string(),
+ );
+ groups
+ .entry(key)
+ .or_default()
+ .push((position.row_position(), position.score()));
+ }
+
+ let mut out = Vec::with_capacity(groups.len());
+ for (key, mut hits) in groups {
+ let source = available
+ .get(&key)
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: format!(
+ "primary-key hybrid result references unknown data file
{}",
+ key.2
+ ),
+ source: None,
+ })?;
+ // Fused positions are physically unique, so no duplicate position
within a
+ // file; sort ascending and coalesce into inclusive ranges with aligned
+ // fused scores.
+ hits.sort_by_key(|(pos, _)| *pos);
+ let mut row_ranges: Vec<RowRange> = Vec::new();
+ let mut scores: Vec<f32> = Vec::with_capacity(hits.len());
+ let mut start = hits[0].0;
+ let mut end = hits[0].0;
+ scores.push(hits[0].1);
+ for &(pos, score) in &hits[1..] {
+ if pos == end + 1 {
+ end = pos;
+ } else {
+ row_ranges.push(RowRange::new(start, end));
+ start = pos;
+ end = pos;
+ }
+ scores.push(score);
+ }
+ row_ranges.push(RowRange::new(start, end));
+
+ out.push(PkVectorIndexedSplit {
+ split: source.clone(),
+ row_ranges,
+ scores: Some(scores),
+ });
+ }
+ Ok(out)
}
#[cfg(feature = "fulltext")]
@@ -502,4 +1095,758 @@ mod tests {
assert_eq!(ranked.row_ids[0], 2);
assert!(ranked.scores[0] > ranked.scores[1]);
}
+
+ // (b) Snapshot pinning: routes that pinned different snapshots must fail
loud
+ // (mirror Java `rankPhysical`'s mismatch guard); equal snapshots are
accepted.
+ #[test]
+ fn check_single_snapshot_rejects_mismatch() {
+ assert_eq!(check_single_snapshot(std::iter::empty()).unwrap(), None);
+ assert_eq!(check_single_snapshot([7, 7, 7]).unwrap(), Some(7));
+ let err = check_single_snapshot([7, 8]).unwrap_err();
+ assert!(
+ format!("{err:?}").contains("same snapshot"),
+ "snapshot mismatch must fail loud, got: {err:?}"
+ );
+ }
+}
+
+#[cfg(all(test, feature = "fulltext"))]
+mod pk_hybrid_tests {
+ use super::*;
+ use crate::catalog::Identifier;
+ use crate::io::{FileIO, FileIOBuilder};
+ use crate::spec::{
+ DataFileMeta, DataType, FloatType, GlobalIndexMeta, IndexFileMeta,
IntType, Schema,
+ TableSchema, VarCharType, VectorType,
+ };
+ use crate::table::pk_full_text_bucket_state::PK_FULL_TEXT_INDEX_TYPE;
+ use crate::table::pk_vector_position_read::{PKEY_VECTOR_POSITION_COLUMN,
SEARCH_SCORE_COLUMN};
+ use crate::table::schema_manager::SchemaManager;
+ use crate::table::{CommitMessage, Table, TableCommit};
+ use arrow_array::{
+ builder::{FixedSizeListBuilder, Float32Builder},
+ Array, ArrayRef, Float32Array, Int32Array, RecordBatch, StringArray,
+ };
+ use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema
as ArrowSchema};
+ use bytes::Bytes;
+ use paimon_ftindex_core::io::PosWriter as FtPosWriter;
+ use paimon_ftindex_core::{FullTextIndexConfig, FullTextIndexWriter};
+ use paimon_vindex_core::index::{VectorIndexConfig, VectorIndexTrainer,
VectorIndexWriter};
+ use paimon_vindex_core::io::PosWriter as VindexPosWriter;
+ use std::collections::HashMap as StdHashMap;
+ use std::sync::Arc;
+
+ const DIM: usize = 4;
+ const VECTOR_COLUMN: &str = "embedding";
+ const TEXT_COLUMN: &str = "body";
+ const VECTOR_INDEX_TYPE: &str = "ivf-flat";
+
+ /// Table options routing BOTH a vector column and a text column to the
+ /// primary-key physical read paths (data-evolution left off).
+ 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(),
+ ),
+ (
+ format!("fields.{VECTOR_COLUMN}.pk-vector.index.type"),
+ VECTOR_INDEX_TYPE.to_string(),
+ ),
+ (
+ format!("fields.{VECTOR_COLUMN}.pk-vector.distance.metric"),
+ "l2".to_string(),
+ ),
+ (
+ "pk-full-text.index.columns".to_string(),
+ TEXT_COLUMN.to_string(),
+ ),
+ ]
+ }
+
+ fn pk_schema(extra: &[(&str, &str)]) -> TableSchema {
+ let mut builder = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(
+ VECTOR_COLUMN,
+ DataType::Vector(
+ VectorType::try_new(true, DIM as u32,
DataType::Float(FloatType::new()))
+ .unwrap(),
+ ),
+ )
+ .column(TEXT_COLUMN, DataType::VarChar(VarCharType::string_type()))
+ .primary_key(["id"]);
+ for (k, v) in table_options() {
+ builder = builder.option(k, v);
+ }
+ for (k, v) in extra {
+ builder = builder.option(*k, *v);
+ }
+ TableSchema::new(0, &builder.build().unwrap())
+ }
+
+ fn data_batch(ids: &[i32], vectors: &[[f32; DIM]], texts: &[&str]) ->
RecordBatch {
+ let element_field = Arc::new(ArrowField::new("element",
ArrowDataType::Float32, true));
+ let mut vector_builder =
FixedSizeListBuilder::new(Float32Builder::new(), DIM as i32)
+ .with_field(element_field.clone());
+ for vector in vectors {
+ for &value in vector {
+ vector_builder.values().append_value(value);
+ }
+ vector_builder.append(true);
+ }
+ let schema = Arc::new(ArrowSchema::new(vec![
+ ArrowField::new("id", ArrowDataType::Int32, false),
+ ArrowField::new(
+ VECTOR_COLUMN,
+ ArrowDataType::FixedSizeList(element_field, DIM as i32),
+ true,
+ ),
+ ArrowField::new(TEXT_COLUMN, ArrowDataType::Utf8, false),
+ ]));
+ RecordBatch::try_new(
+ schema,
+ vec![
+ Arc::new(Int32Array::from(ids.to_vec())) as ArrayRef,
+ Arc::new(vector_builder.finish()) as ArrayRef,
+ Arc::new(StringArray::from(texts.to_vec())) as ArrayRef,
+ ],
+ )
+ .unwrap()
+ }
+
+ fn java_write_utf(s: &str) -> Vec<u8> {
+ let mut body = Vec::new();
+ for c in s.encode_utf16() {
+ if (0x0001..=0x007F).contains(&c) {
+ body.push(c as u8);
+ } else if c > 0x07FF {
+ body.push(0xE0 | (c >> 12) as u8);
+ body.push(0x80 | ((c >> 6) & 0x3F) as u8);
+ body.push(0x80 | (c & 0x3F) as u8);
+ } else {
+ body.push(0xC0 | (c >> 6) as u8);
+ body.push(0x80 | (c & 0x3F) as u8);
+ }
+ }
+ let mut out = (body.len() as u16).to_be_bytes().to_vec();
+ out.extend_from_slice(&body);
+ out
+ }
+
+ fn source_meta_bytes(data_level: i32, files: &[(&str, i64)]) -> Vec<u8> {
+ let mut out = Vec::new();
+ out.extend_from_slice(&1i32.to_be_bytes());
+ out.extend_from_slice(&data_level.to_be_bytes());
+ out.extend_from_slice(&(files.len() as i32).to_be_bytes());
+ for (name, rows) in files {
+ out.extend_from_slice(&java_write_utf(name));
+ out.extend_from_slice(&rows.to_be_bytes());
+ }
+ out
+ }
+
+ async fn write_bytes(file_io: &FileIO, path: &str, bytes: Vec<u8>) {
+ file_io
+ .new_output(path)
+ .unwrap()
+ .write(Bytes::from(bytes))
+ .await
+ .unwrap();
+ }
+
+ /// Build a real vindex IVF-flat ANN segment over `vectors` (label ==
physical
+ /// position); `nlist = 1` keeps the search exact.
+ async fn write_ann_segment(
+ file_io: &FileIO,
+ location: &str,
+ file_name: &str,
+ vectors: &[[f32; DIM]],
+ ) -> u64 {
+ let n = vectors.len();
+ let flat: Vec<f32> = vectors.iter().flat_map(|v|
v.iter().copied()).collect();
+ let ids: Vec<i64> = (0..n as i64).collect();
+ let native_options = StdHashMap::from([
+ ("index.type".to_string(), "ivf_flat".to_string()),
+ ("dimension".to_string(), DIM.to_string()),
+ ("nlist".to_string(), "1".to_string()),
+ ("metric".to_string(), "l2".to_string()),
+ ]);
+ let config = VectorIndexConfig::from_options(&native_options).unwrap();
+ let training = VectorIndexTrainer::train(config, &flat, n).unwrap();
+ let mut writer = VectorIndexWriter::new(training);
+ writer.add_vectors(&ids, &flat, n).unwrap();
+ let mut bytes = Vec::new();
+ {
+ let mut output = VindexPosWriter::new(&mut bytes);
+ writer.write(&mut output).unwrap();
+ }
+ let file_size = bytes.len() as u64;
+ write_bytes(file_io, &format!("{location}/index/{file_name}"),
bytes).await;
+ file_size
+ }
+
+ /// Build a full-text archive (row id == physical position) via the native
core.
+ fn build_archive(docs: &[(i64, &str)]) -> Vec<u8> {
+ let mut writer =
FullTextIndexWriter::new(FullTextIndexConfig::new()).unwrap();
+ for (row_id, text) in docs {
+ writer.add_document(*row_id, (*text).to_string()).unwrap();
+ }
+ let mut out = FtPosWriter::new(Vec::<u8>::new());
+ writer.write(&mut out).unwrap();
+ out.into_inner()
+ }
+
+ async fn 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", "pk_hybrid"),
+ location.to_string(),
+ (*schema).clone(),
+ None,
+ )
+ }
+
+ /// Build a complete self-contained primary-key hybrid table on an
in-memory
+ /// FileIO: persist the schema, write a real data file, build+commit a real
+ /// vindex ANN segment AND a full-text archive over the SAME compacted data
+ /// file. Returns the opened table.
+ async fn build_hybrid_table(
+ location: &str,
+ ids: &[i32],
+ vectors: &[[f32; DIM]],
+ texts: &[&str],
+ extra_options: &[(&str, &str)],
+ ) -> Table {
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ for dir in ["schema", "snapshot", "manifest", "index"] {
+ file_io.mkdirs(&format!("{location}/{dir}")).await.unwrap();
+ }
+ let schema = pk_schema(extra_options);
+ write_bytes(
+ &file_io,
+ &format!("{location}/schema/schema-{}", schema.id()),
+ serde_json::to_vec(&schema).unwrap(),
+ )
+ .await;
+
+ let table = open_table(&file_io, location).await;
+
+ // Write a real data file via the public write path.
+ let write_builder = table.new_write_builder();
+ let mut writer = write_builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&data_batch(ids, vectors, texts))
+ .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();
+ let data_file_name = base_meta.file_name.clone();
+ let row_count = base_meta.row_count;
+
+ // Only a compacted, non-level-0 file backs the primary-key indices.
+ let indexed_meta = DataFileMeta {
+ level: 1,
+ file_source: Some(1),
+ first_row_id: Some(0),
+ ..base_meta
+ };
+
+ let vector_field_id = schema
+ .fields()
+ .iter()
+ .find(|f| f.name() == VECTOR_COLUMN)
+ .unwrap()
+ .id();
+ let text_field_id = schema
+ .fields()
+ .iter()
+ .find(|f| f.name() == TEXT_COLUMN)
+ .unwrap()
+ .id();
+
+ // Vector ANN segment.
+ let vector_index_name = "vector-ivf-flat-pk-hybrid.index".to_string();
+ let vector_index_size =
+ write_ann_segment(&file_io, location, &vector_index_name,
vectors).await;
+ let vector_index = IndexFileMeta {
+ index_type: VECTOR_INDEX_TYPE.to_string(),
+ file_name: vector_index_name,
+ file_size: i64::try_from(vector_index_size).unwrap(),
+ row_count: i32::try_from(row_count).unwrap(),
+ deletion_vectors_ranges: None,
+ global_index_meta: Some(GlobalIndexMeta {
+ row_range_start: 0,
+ row_range_end: row_count - 1,
+ index_field_id: vector_field_id,
+ extra_field_ids: None,
+ source_meta: Some(source_meta_bytes(
+ indexed_meta.level,
+ &[(&data_file_name, row_count)],
+ )),
+ index_meta: None,
+ }),
+ };
+
+ // Full-text archive (row id == physical position).
+ let ft_index_name = "full-text-pk-hybrid.index".to_string();
+ let docs: Vec<(i64, &str)> = texts
+ .iter()
+ .enumerate()
+ .map(|(pos, text)| (pos as i64, *text))
+ .collect();
+ let archive = build_archive(&docs);
+ write_bytes(
+ &file_io,
+ &format!("{location}/index/{ft_index_name}"),
+ archive,
+ )
+ .await;
+ let ft_index = IndexFileMeta {
+ index_type: PK_FULL_TEXT_INDEX_TYPE.to_string(),
+ file_name: ft_index_name,
+ file_size: 1,
+ row_count: i32::try_from(row_count).unwrap(),
+ deletion_vectors_ranges: None,
+ global_index_meta: Some(GlobalIndexMeta {
+ row_range_start: 0,
+ row_range_end: row_count - 1,
+ index_field_id: text_field_id,
+ extra_field_ids: None,
+ source_meta: Some(source_meta_bytes(1, &[(&data_file_name,
row_count)])),
+ index_meta: None,
+ }),
+ };
+
+ let mut message = CommitMessage::new(partition, bucket,
vec![indexed_meta]);
+ message.new_index_files = vec![vector_index, ft_index];
+ TableCommit::new(table.clone(), "pk-hybrid".to_string())
+ .commit(vec![message])
+ .await
+ .unwrap();
+
+ table
+ }
+
+ fn column_i32(batches: &[RecordBatch], name: &str) -> Vec<i32> {
+ batches
+ .iter()
+ .flat_map(|b| {
+ let idx = b.schema().index_of(name).unwrap();
+ b.column(idx)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap()
+ .values()
+ .to_vec()
+ })
+ .collect()
+ }
+
+ fn column_f32(batches: &[RecordBatch], name: &str) -> Vec<f32> {
+ batches
+ .iter()
+ .flat_map(|b| {
+ let idx = b.schema().index_of(name).unwrap();
+ b.column(idx)
+ .as_any()
+ .downcast_ref::<Float32Array>()
+ .unwrap()
+ .values()
+ .to_vec()
+ })
+ .collect()
+ }
+
+ // (a) End-to-end PK hybrid: fuse a vector route and a full-text route,
best
+ // fused first, with the score column and without the internal position
column.
+ // Data is chosen so the fused order differs from BOTH route-only orders,
+ // proving the ranker actually combines the two routes.
+ #[tokio::test]
+ async fn end_to_end_pk_hybrid_fuses_vector_and_full_text() {
+ // Vector nearest -> pos0,1,2,3 (query [10,0,0,0]).
+ let vectors = vec![
+ [10.0, 0.0, 0.0, 0.0], // pos0 dist 0
+ [9.0, 0.0, 0.0, 0.0], // pos1 dist 1
+ [8.0, 0.0, 0.0, 0.0], // pos2 dist 4
+ [7.0, 0.0, 0.0, 0.0], // pos3 dist 9
+ ];
+ // Full-text "alpha" -> pos2 (tf 3) then pos0 (tf 1).
+ let texts = vec!["alpha", "beta", "alpha alpha alpha", "gamma"];
+ let ids = [100, 101, 102, 103];
+ let table = build_hybrid_table("memory:/pk_hybrid_e2e", &ids,
&vectors, &texts, &[]).await;
+
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![10.0, 0.0, 0.0, 0.0],
+ 4,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap()
+ .add_full_text_route(
+ TEXT_COLUMN,
+ r#"{"match":{"query":"alpha"}}"#,
+ 4,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(3).with_rrf_ranker();
+
+ let batches = builder
+ .execute_read()
+ .await
+ .expect("pk hybrid execute_read failed")
+ .try_collect::<Vec<_>>()
+ .await
+ .expect("collecting hybrid read batches failed");
+
+ // RRF: pos0 = 1/61 + 1/62 (vec r1, ft r2); pos2 = 1/63 + 1/61 (vec
r3, ft r1);
+ // pos1 = 1/62 (vec r2). Fused best-first: 100, 102, 101. This differs
from
+ // the vector-only order (100,101,102) and the full-text-only order
(102,100),
+ // so it can only come from fusing BOTH routes.
+ let ids_out = column_i32(&batches, "id");
+ assert_eq!(ids_out, vec![100, 102, 101], "fused best-first order");
+
+ // The unified score column is present, descending (best first); the
internal
+ // position column is stripped.
+ let scores = column_f32(&batches, SEARCH_SCORE_COLUMN);
+ assert_eq!(scores.len(), 3);
+ assert!(
+ scores[0] >= scores[1] && scores[1] >= scores[2],
+ "scores must be best-first: {scores:?}"
+ );
+ for batch in &batches {
+ assert!(
+ batch
+ .schema()
+ .index_of(PKEY_VECTOR_POSITION_COLUMN)
+ .is_err(),
+ "internal position column must be stripped"
+ );
+ assert!(
+ batch.schema().index_of(SEARCH_SCORE_COLUMN).is_ok(),
+ "score column must be present"
+ );
+ }
+ }
+
+ // (c) A mixed PK/global route set must fail loud on execute_read.
+ #[tokio::test]
+ async fn mixed_pk_and_global_routes_fail_loud() {
+ let table = build_hybrid_table(
+ "memory:/pk_hybrid_mixed",
+ &[100, 101],
+ &[[10.0, 0.0, 0.0, 0.0], [9.0, 0.0, 0.0, 0.0]],
+ &["alpha", "beta"],
+ &[],
+ )
+ .await;
+
+ // The vector route is a PK route; the full-text route targets a
column NOT
+ // configured for PK full-text -> global address space -> mixed ->
fail loud.
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![10.0, 0.0, 0.0, 0.0],
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap()
+ .add_full_text_route(
+ "not_indexed",
+ r#"{"match":{"query":"alpha"}}"#,
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(2);
+
+ let err = match builder.execute_read().await {
+ Ok(_) => panic!("mixed PK/global routes must fail loud"),
+ Err(e) => e,
+ };
+ assert!(
+ format!("{err:?}").contains("address spaces"),
+ "mixed PK/global must fail loud, got: {err:?}"
+ );
+ }
+
+ // (d) The PK hybrid path produces physical positions, not global row ids:
+ // execute / execute_scored must fail loud and point at execute_read.
+ #[tokio::test]
+ async fn pk_hybrid_execute_and_scored_fail_loud() {
+ let table = build_hybrid_table(
+ "memory:/pk_hybrid_guard",
+ &[100, 101],
+ &[[10.0, 0.0, 0.0, 0.0], [9.0, 0.0, 0.0, 0.0]],
+ &["alpha", "beta"],
+ &[],
+ )
+ .await;
+
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![10.0, 0.0, 0.0, 0.0],
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap()
+ .add_full_text_route(
+ TEXT_COLUMN,
+ r#"{"match":{"query":"alpha"}}"#,
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(2);
+
+ let scored_err = builder.execute_scored().await.unwrap_err();
+ assert!(
+ format!("{scored_err:?}").contains("execute_read"),
+ "PK hybrid execute_scored must point at execute_read, got:
{scored_err:?}"
+ );
+ let execute_err = builder.execute().await.unwrap_err();
+ assert!(
+ format!("{execute_err:?}").contains("execute_read"),
+ "PK hybrid execute must point at execute_read, got:
{execute_err:?}"
+ );
+ }
+
+ /// A primary-key hybrid table whose user schema carries an extra column
named
+ /// `reserved`, used to prove the materialized read rejects a reserved
metadata
+ /// name arriving via the default (all-columns) projection.
+ fn pk_hybrid_table_with_reserved_column(reserved: &str) -> Table {
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let mut builder = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(
+ VECTOR_COLUMN,
+ DataType::Vector(
+ VectorType::try_new(true, DIM as u32,
DataType::Float(FloatType::new()))
+ .unwrap(),
+ ),
+ )
+ .column(TEXT_COLUMN, DataType::VarChar(VarCharType::string_type()))
+ .column(reserved, DataType::VarChar(VarCharType::string_type()))
+ .primary_key(["id"]);
+ for (k, v) in table_options() {
+ builder = builder.option(k, v);
+ }
+ let schema = TableSchema::new(0, &builder.build().unwrap());
+ Table::new(
+ file_io,
+ Identifier::new("default", "pk_hybrid_reserved"),
+ "memory:/pk_hybrid_reserved".to_string(),
+ schema,
+ None,
+ )
+ }
+
+ // A user column colliding with an injected metadata column
+ // (`_PKEY_VECTOR_POSITION`, `__paimon_search_score`, or `_ROW_ID`) must
make the
+ // primary-key hybrid materialized read fail loud up front — before any
route
+ // runs and regardless of results — mirroring the primary-key vector read
guard.
+ #[tokio::test]
+ async fn pk_hybrid_read_rejects_reserved_user_column() {
+ for reserved in ["_PKEY_VECTOR_POSITION", "__paimon_search_score",
"_ROW_ID"] {
+ let table = pk_hybrid_table_with_reserved_column(reserved);
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![1.0, 0.0, 0.0, 0.0],
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(2);
+ let err = builder
+ .execute_read()
+ .await
+ .err()
+ .expect("reserved user column must fail loud");
+ assert!(
+ matches!(&err, crate::Error::DataInvalid { message, .. }
+ if message.contains("reserved column")),
+ "unexpected error for {reserved}: {err:?}"
+ );
+ }
+ }
+
+ // (e) DE regression: a hybrid over non-PK routes (no PK configs) still
takes the
+ // global-row-id execute_scored path unchanged. An empty table yields an
empty
+ // result without hitting the PK guard.
+ #[tokio::test]
+ async fn global_hybrid_execute_scored_unchanged() {
+ // A table without any PK index configuration: both routes resolve to
the
+ // append/data-evolution global path.
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(VECTOR_COLUMN, {
+ DataType::Vector(
+ VectorType::try_new(true, DIM as u32,
DataType::Float(FloatType::new()))
+ .unwrap(),
+ )
+ })
+ .column(TEXT_COLUMN,
DataType::VarChar(VarCharType::string_type()))
+ .option("bucket", "1")
+ .option("row-tracking.enabled", "true")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "global_hybrid"),
+ "memory:/global_hybrid".to_string(),
+ schema,
+ None,
+ );
+
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![1.0, 0.0, 0.0, 0.0],
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap()
+ .add_full_text_route(
+ TEXT_COLUMN,
+ r#"{"match":{"query":"alpha"}}"#,
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(2);
+
+ // No snapshot -> both routes empty -> fused empty; the DE path must
NOT hit
+ // the PK fail-loud guard.
+ let result = builder
+ .execute_scored()
+ .await
+ .expect("DE hybrid path failed");
+ assert!(
+ result.is_empty(),
+ "empty table yields empty DE hybrid result"
+ );
+ }
+
+ // (f) FAST-only guard: a primary-key hybrid with a full-text route under a
+ // non-FAST global-index search mode (FULL) must fail loud, mirroring
+ // `full_text_search_builder::execute_read` and Java
`PrimaryKeyFullTextRead`.
+ #[tokio::test]
+ async fn pk_hybrid_full_text_route_rejects_non_fast_mode() {
+ let table = build_hybrid_table(
+ "memory:/pk_hybrid_full_mode",
+ &[100, 101, 102, 103],
+ &[
+ [10.0, 0.0, 0.0, 0.0],
+ [9.0, 0.0, 0.0, 0.0],
+ [8.0, 0.0, 0.0, 0.0],
+ [7.0, 0.0, 0.0, 0.0],
+ ],
+ &["alpha", "beta", "alpha alpha alpha", "gamma"],
+ &[("global-index.search-mode", "full")],
+ )
+ .await;
+
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![10.0, 0.0, 0.0, 0.0],
+ 4,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap()
+ .add_full_text_route(
+ TEXT_COLUMN,
+ r#"{"match":{"query":"alpha"}}"#,
+ 4,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(3);
+
+ let err = match builder.execute_read().await {
+ Ok(_) => panic!("non-FAST global-index search mode must fail
loud"),
+ Err(e) => e,
+ };
+ assert!(
+ format!("{err:?}").contains("FAST"),
+ "PK hybrid full-text route under FULL mode must fail loud, got:
{err:?}"
+ );
+ }
+
+ // (g) Snapshot pinning: the read-latest path resolves ONE snapshot up
front and
+ // passes it to every route (mirror Java `routeBuilders()`). The pinned
table
+ // carries `scan.snapshot-id` set to the resolved latest snapshot, so both
+ // routes plan against the same version instead of each re-resolving
latest.
+ #[tokio::test]
+ async fn pk_hybrid_read_latest_pins_one_snapshot_for_all_routes() {
+ let table = build_hybrid_table(
+ "memory:/pk_hybrid_pin",
+ &[100, 101],
+ &[[10.0, 0.0, 0.0, 0.0], [9.0, 0.0, 0.0, 0.0]],
+ &["alpha", "beta"],
+ &[],
+ )
+ .await;
+
+ let mut builder = table.new_hybrid_search_builder();
+ builder
+ .add_vector_route(
+ VECTOR_COLUMN,
+ vec![10.0, 0.0, 0.0, 0.0],
+ 2,
+ 1.0,
+ HashMap::new(),
+ )
+ .unwrap();
+ builder.with_limit(2);
+
+ let pinned = builder
+ .resolve_pinned_route_table()
+ .await
+ .expect("pinning must succeed")
+ .expect("read-latest path must pin a snapshot");
+ // First commit -> snapshot 1; the pinned copy targets exactly it.
+ assert_eq!(
+ pinned.schema().options().get(SCAN_SNAPSHOT_ID_OPTION),
+ Some(&"1".to_string()),
+ "read-latest hybrid must pin the resolved latest snapshot id"
+ );
+ }
}
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index ca6db9c8..be680cb0 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -65,6 +65,8 @@ mod pk_full_text_bucket_state;
mod pk_full_text_read;
#[cfg(feature = "fulltext")]
mod pk_full_text_scan;
+mod pk_search_position;
+mod pk_search_ranker;
mod pk_vector_data_file_reader;
mod pk_vector_indexed_split_read;
mod pk_vector_orchestrator;
diff --git a/crates/paimon/src/table/pk_full_text_read.rs
b/crates/paimon/src/table/pk_full_text_read.rs
index a8dd70a8..3c5a34b1 100644
--- a/crates/paimon/src/table/pk_full_text_read.rs
+++ b/crates/paimon/src/table/pk_full_text_read.rs
@@ -130,7 +130,7 @@ pub(crate) fn top_k_by_score(
/// scores are already final relevance scores, so any transform would corrupt
them
/// (spec D1). Groups are emitted in ascending group-key order for a
deterministic
/// file/position materialization order; the caller reorders back to
best-first.
-fn build_full_text_indexed_splits(
+pub(crate) fn build_full_text_indexed_splits(
survivors: Vec<PrimaryKeyFullTextCandidate>,
splits: &[PrimaryKeyFullTextSearchSplit],
) -> crate::Result<Vec<PkVectorIndexedSplit>> {
@@ -246,6 +246,21 @@ fn build_full_text_indexed_splits(
Ok(out)
}
+/// The primary-key full-text route's search output plus the source context a
+/// later materialization (or a hybrid fusion across routes) needs.
`candidates`
+/// are the best-score-first survivors; `splits` are the planned per-bucket
source
+/// splits their `split_index` refers into; `snapshot_id` is the snapshot the
plan
+/// resolved. The splits are borrowed from the originating plan (kept alive by
the
+/// caller) because `split_index` is only meaningful against that plan.
Produced by
+/// [`PrimaryKeyFullTextRead::search_route`].
+pub(crate) struct PrimaryKeyFullTextRouteResult<'a> {
+ pub(crate) candidates: Vec<PrimaryKeyFullTextCandidate>,
+ // Read by the hybrid route consumer (fuses routes before materializing);
the
+ // materialized read reaches candidates directly and holds the plan itself.
+ pub(crate) snapshot_id: i64,
+ pub(crate) splits: &'a [PrimaryKeyFullTextSearchSplit],
+}
+
/// FAST-only primary-key full-text materialized read. Given a planned set of
/// per-bucket search splits, it searches each bucket through the full-text
archive
/// reader, fuses the hits cross-bucket by score, materializes the winning
physical
@@ -277,16 +292,19 @@ impl PrimaryKeyFullTextRead {
}
}
- /// Search every planned bucket, fuse the hits by score into a global
Top-`limit`,
- /// materialize the winning rows, and emit them best-score-first with the
- /// unified score column. An empty plan or an empty result yields an empty
- /// stream. `query` is passed verbatim to the archive reader (spec D2).
- pub(crate) async fn read(
+ /// Search every planned bucket and fuse the hits by score into the global
+ /// best-score-first Top-`limit`, WITHOUT materializing any rows. Returns
the
+ /// surviving candidates together with the route source context (the plan's
+ /// resolved snapshot id and its per-bucket source splits) so a caller can
fuse
+ /// them with another route before materializing. [`read`](Self::read) is
+ /// layered on top of this. `query` is passed verbatim to the archive
reader
+ /// (spec D2).
+ pub(crate) async fn search_route<'p>(
&self,
- plan: &PrimaryKeyFullTextScanPlan,
+ plan: &'p PrimaryKeyFullTextScanPlan,
query: &str,
limit: usize,
- ) -> crate::Result<ArrowRecordBatchStream> {
+ ) -> crate::Result<PrimaryKeyFullTextRouteResult<'p>> {
// Every planned split must belong to the plan's resolved snapshot;
mixing
// snapshots would search/materialize physical rows against the wrong
// version. The scan already pins one snapshot, so a mismatch is
malformed
@@ -334,6 +352,24 @@ impl PrimaryKeyFullTextRead {
// Global cross-bucket fusion: score-descending Top-`limit`.
let survivors = top_k_by_score(candidates, limit);
+ Ok(PrimaryKeyFullTextRouteResult {
+ candidates: survivors,
+ snapshot_id: plan.snapshot_id,
+ splits: &plan.splits,
+ })
+ }
+
+ /// Search every planned bucket, fuse the hits by score into a global
Top-`limit`,
+ /// materialize the winning rows, and emit them best-score-first with the
+ /// unified score column. An empty plan or an empty result yields an empty
+ /// stream. `query` is passed verbatim to the archive reader (spec D2).
+ pub(crate) async fn read(
+ &self,
+ plan: &PrimaryKeyFullTextScanPlan,
+ query: &str,
+ limit: usize,
+ ) -> crate::Result<ArrowRecordBatchStream> {
+ let survivors = self.search_route(plan, query,
limit).await?.candidates;
if survivors.is_empty() {
return Ok(Box::pin(stream::empty()));
}
@@ -966,4 +1002,90 @@ mod read_tests {
.await;
assert!(batches.is_empty(), "no match must yield an empty stream");
}
+
+ // ---- search_route: candidate-only producer returns candidates + context
----
+ #[tokio::test]
+ async fn search_route_returns_candidates_and_source_context() {
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let table_path = "memory:/pk_ft_route";
+ // pos2 ("alpha alpha alpha") scores higher for "alpha" than pos0
("alpha").
+ let archive = build_archive(&[
+ (0, "alpha"),
+ (1, "beta"),
+ (2, "alpha alpha alpha"),
+ (3, "gamma"),
+ ]);
+ write_bytes(&file_io, &format!("{table_path}/index/ft-0"),
archive).await;
+
+ let (reader, split) = build_data(&file_io, table_path, vec![100, 101,
102, 103], &[]).await;
+ let plan = PrimaryKeyFullTextScanPlan {
+ snapshot_id: 1,
+ splits: vec![PrimaryKeyFullTextSearchSplit::new(
+ split,
+ vec![ft_payload("ft-0", &[("d0.mosaic", 4)])],
+ Vec::new(),
+ )
+ .unwrap()],
+ };
+
+ let read = PrimaryKeyFullTextRead::new(file_io.clone(), reader,
table_path.to_string());
+ let route = read
+ .search_route(&plan, r#"{"match":{"query":"alpha"}}"#, 10)
+ .await
+ .unwrap();
+
+ // Two docs contain "alpha" -> two best-score-first candidates, no
+ // materialization performed.
+ assert_eq!(route.candidates.len(), 2, "both alpha hits must survive");
+ assert!(
+ route.candidates[0].score >= route.candidates[1].score,
+ "candidates must be best-score-first: {:?}",
+ route.candidates.iter().map(|c| c.score).collect::<Vec<_>>()
+ );
+ // The strong "alpha alpha alpha" hit (physical position 2) ranks
first.
+ assert_eq!(route.candidates[0].row_position, 2);
+ assert_eq!(route.candidates[1].row_position, 0);
+
+ // Source context: the plan's snapshot and its per-bucket source
splits are
+ // carried through so a caller can materialize (or fuse) without the
plan.
+ assert_eq!(route.snapshot_id, 1);
+ assert_eq!(route.splits.len(), 1, "one bucket split expected");
+ // The candidate's split_index refers into the returned splits.
+ assert!(route
+ .candidates
+ .iter()
+ .all(|c| c.split_index < route.splits.len()));
+ }
+
+ /// A real (non-empty) plan whose query matches nothing still reports the
+ /// plan's pinned snapshot id and its source splits — a zero-candidate
route
+ /// must not lose the snapshot the cross-route consistency guard requires.
+ #[tokio::test]
+ async fn search_route_zero_candidates_still_carries_snapshot() {
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let table_path = "memory:/pk_ft_route_empty";
+ let archive = build_archive(&[(0, "alpha"), (1, "beta")]);
+ write_bytes(&file_io, &format!("{table_path}/index/ft-0"),
archive).await;
+
+ let (reader, split) = build_data(&file_io, table_path, vec![100, 101],
&[]).await;
+ let plan = PrimaryKeyFullTextScanPlan {
+ snapshot_id: 1,
+ splits: vec![PrimaryKeyFullTextSearchSplit::new(
+ split,
+ vec![ft_payload("ft-0", &[("d0.mosaic", 2)])],
+ Vec::new(),
+ )
+ .unwrap()],
+ };
+
+ let read = PrimaryKeyFullTextRead::new(file_io.clone(), reader,
table_path.to_string());
+ let route = read
+ .search_route(&plan, r#"{"match":{"query":"zeta"}}"#, 10)
+ .await
+ .unwrap();
+
+ assert!(route.candidates.is_empty(), "no doc matches 'zeta'");
+ assert_eq!(route.snapshot_id, 1, "real snapshot id survives zero
hits");
+ assert_eq!(route.splits.len(), 1, "source splits still carried
through");
+ }
}
diff --git a/crates/paimon/src/table/pk_full_text_scan.rs
b/crates/paimon/src/table/pk_full_text_scan.rs
index 081ece10..f6093858 100644
--- a/crates/paimon/src/table/pk_full_text_scan.rs
+++ b/crates/paimon/src/table/pk_full_text_scan.rs
@@ -249,8 +249,10 @@ impl BucketAccumulator {
/// The per-bucket search splits produced by planning.
pub(crate) struct PrimaryKeyFullTextScanPlan {
- // The snapshot the plan was resolved against; the read guards every split
- // against it before searching.
+ // The snapshot the plan resolved during planning (pinned before the index
+ // manifest is read); the read guards every split against it before
searching.
+ // It is authoritative even when planning yields zero splits, and is `0`
only
+ // for a table with no snapshot at all (never written).
pub snapshot_id: i64,
pub splits: Vec<PrimaryKeyFullTextSearchSplit>,
}
@@ -283,21 +285,28 @@ impl<'a> PrimaryKeyFullTextScan<'a> {
if let Some(filter) = &self.filter {
read_builder.with_filter(filter.clone());
}
- let data_splits = read_builder
+ // Plan the data splits and capture the snapshot the scan pinned in one
+ // pass. The trace carries the resolved snapshot id even when the scan
+ // yields zero data splits, so the plan reports its real snapshot id
+ // (required by the cross-route snapshot-consistency guard) instead of
+ // deriving it from a first split that may not exist.
+ let (data_plan, trace) = read_builder
.new_scan()
.with_scan_all_files()
- .plan()
- .await?
- .splits()
- .to_vec();
-
- let Some(first_split) = data_splits.first() else {
+ .plan_with_trace()
+ .await?;
+ let data_splits = data_plan.splits().to_vec();
+
+ // No snapshot at all (table never written): nothing to search and no
+ // snapshot to pin. The empty split list makes every downstream
consumer
+ // treat this as "no candidates", so this is the only plan without a
real
+ // snapshot id.
+ let Some(snapshot_id) = trace.snapshot_id else {
return Ok(PrimaryKeyFullTextScanPlan {
snapshot_id: 0,
splits: Vec::new(),
});
};
- let snapshot_id = first_split.snapshot_id();
let snapshot = snapshot_manager.get_snapshot(snapshot_id).await?;
let mut entries: Vec<(BinaryRow, i32, IndexFileMeta)> = Vec::new();
diff --git a/crates/paimon/src/table/pk_search_position.rs
b/crates/paimon/src/table/pk_search_position.rs
new file mode 100644
index 00000000..5cf5f3e2
--- /dev/null
+++ b/crates/paimon/src/table/pk_search_position.rs
@@ -0,0 +1,250 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Shared scored physical row position for primary-key search.
+//!
+//! Rust mirror of Java
+//! `org.apache.paimon.table.source.PrimaryKeySearchPosition`: a physical row
+//! address `(partition, bucket, data_file_name, row_position)` carrying a
search
+//! `score`. Equality and hashing are on the physical identity ONLY (the score
is
+//! excluded) so hits fused from different search routes (vector, full-text)
that
+//! resolve to the same physical row collapse to one position regardless of
their
+//! per-route scores.
+
+use std::hash::{Hash, Hasher};
+
+use crate::spec::BinaryRow;
+
+fn data_invalid(message: impl Into<String>) -> crate::Error {
+ crate::Error::DataInvalid {
+ message: message.into(),
+ source: None,
+ }
+}
+
+/// A scored physical row position in a primary-key table snapshot. Mirrors
Java
+/// `PrimaryKeySearchPosition`.
+#[derive(Clone, Debug)]
+pub(crate) struct PrimaryKeySearchPosition {
+ partition: BinaryRow,
+ bucket: i32,
+ data_file_name: String,
+ row_position: i64,
+ score: f32,
+}
+
+impl PrimaryKeySearchPosition {
+ /// Build a position, rejecting a negative `row_position` or a non-finite
+ /// `score` (NaN / ±Infinity). A negative position means a bogus physical
+ /// address and a non-finite score would poison any score ordering, so fail
+ /// loud rather than propagate it. Mirrors the `checkArgument`s in Java
+ /// `PrimaryKeySearchPosition`.
+ pub(crate) fn new(
+ partition: BinaryRow,
+ bucket: i32,
+ data_file_name: String,
+ row_position: i64,
+ score: f32,
+ ) -> crate::Result<Self> {
+ if row_position < 0 {
+ return Err(data_invalid(format!(
+ "row position must not be negative, got {row_position} for
{data_file_name}"
+ )));
+ }
+ if !score.is_finite() {
+ return Err(data_invalid(format!(
+ "search score must be finite, got {score} for {data_file_name}
@ {row_position}"
+ )));
+ }
+ Ok(Self {
+ partition,
+ bucket,
+ data_file_name,
+ row_position,
+ score,
+ })
+ }
+
+ pub(crate) fn partition(&self) -> &BinaryRow {
+ &self.partition
+ }
+
+ pub(crate) fn bucket(&self) -> i32 {
+ self.bucket
+ }
+
+ pub(crate) fn data_file_name(&self) -> &str {
+ &self.data_file_name
+ }
+
+ pub(crate) fn row_position(&self) -> i64 {
+ self.row_position
+ }
+
+ pub(crate) fn score(&self) -> f32 {
+ self.score
+ }
+
+ /// Returns a copy of this physical position carrying `new_score`,
re-running
+ /// the finite-score validation. Mirrors Java
`PrimaryKeySearchPosition.withScore`.
+ pub(crate) fn with_score(&self, new_score: f32) -> crate::Result<Self> {
+ Self::new(
+ self.partition.clone(),
+ self.bucket,
+ self.data_file_name.clone(),
+ self.row_position,
+ new_score,
+ )
+ }
+
+ pub(crate) fn from_vector_candidate(
+ candidate: &crate::table::pk_vector_orchestrator::PkVectorCandidate,
+ metric: crate::vindex::pkvector::metric::VectorSearchMetric,
+ ) -> crate::Result<Self> {
+ Self::new(
+ candidate.partition.clone(),
+ candidate.bucket,
+ candidate.data_file_name.clone(),
+ candidate.row_position,
+ metric.distance_to_score(candidate.distance),
+ )
+ }
+
+ #[cfg(feature = "fulltext")]
+ pub(crate) fn from_full_text_candidate(
+ candidate:
&crate::table::pk_full_text_read::PrimaryKeyFullTextCandidate,
+ ) -> crate::Result<Self> {
+ Self::new(
+ candidate.partition.clone(),
+ candidate.bucket,
+ candidate.data_file_name.clone(),
+ candidate.row_position,
+ candidate.score,
+ )
+ }
+}
+
+impl PartialEq for PrimaryKeySearchPosition {
+ /// Physical identity only: `(partition, bucket, data_file_name,
+ /// row_position)`. The `score` is deliberately excluded (mirrors Java
+ /// `PrimaryKeySearchPosition.equals`). The partition is compared by its
+ /// serialized bytes, matching how the vector/full-text routes group
+ /// partitions and keeping equality consistent with [`Hash`].
+ fn eq(&self, other: &Self) -> bool {
+ self.bucket == other.bucket
+ && self.row_position == other.row_position
+ && self.data_file_name == other.data_file_name
+ && self.partition.to_serialized_bytes() ==
other.partition.to_serialized_bytes()
+ }
+}
+
+impl Eq for PrimaryKeySearchPosition {}
+
+impl Hash for PrimaryKeySearchPosition {
+ /// Hashes the same physical-identity fields [`PartialEq`] compares (score
+ /// excluded). `BinaryRow` is not `Hash`, so the partition is hashed via
its
+ /// serialized bytes, exactly the comparator the vector route uses.
+ fn hash<H: Hasher>(&self, state: &mut H) {
+ self.partition.to_serialized_bytes().hash(state);
+ self.bucket.hash(state);
+ self.data_file_name.hash(state);
+ self.row_position.hash(state);
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::table::pk_vector_orchestrator::PkVectorCandidate;
+ use crate::vindex::pkvector::metric::VectorSearchMetric;
+ use std::collections::HashSet;
+
+ fn pos(row_position: i64, score: f32) ->
crate::Result<PrimaryKeySearchPosition> {
+ PrimaryKeySearchPosition::new(BinaryRow::new(0), 0, "f".to_string(),
row_position, score)
+ }
+
+ #[test]
+ fn new_rejects_negative_row_position() {
+ assert!(pos(-1, 1.0).is_err());
+ assert!(pos(0, 1.0).is_ok());
+ }
+
+ #[test]
+ fn new_rejects_non_finite_score() {
+ assert!(pos(0, f32::NAN).is_err());
+ assert!(pos(0, f32::INFINITY).is_err());
+ assert!(pos(0, f32::NEG_INFINITY).is_err());
+ assert!(pos(0, 1.5).is_ok());
+ }
+
+ #[test]
+ fn equal_and_hash_equal_ignoring_score() {
+ let a = pos(3, 0.1).unwrap();
+ let b = pos(3, 0.9).unwrap();
+ assert_eq!(a, b, "positions differing only in score must be equal");
+ let mut set = HashSet::new();
+ set.insert(a);
+ set.insert(b);
+ assert_eq!(set.len(), 1, "score must not affect hash-identity");
+ }
+
+ #[test]
+ fn differing_physical_identity_is_not_equal() {
+ let a = pos(3, 0.5).unwrap();
+ let b = pos(4, 0.5).unwrap();
+ assert_ne!(a, b);
+ let mut set = HashSet::new();
+ set.insert(a);
+ set.insert(b);
+ assert_eq!(set.len(), 2);
+ }
+
+ fn vector_candidate(distance: f32) -> PkVectorCandidate {
+ PkVectorCandidate {
+ split_index: 0,
+ partition: BinaryRow::new(0),
+ bucket: 0,
+ data_file_name: "f".to_string(),
+ row_position: 0,
+ distance,
+ }
+ }
+
+ #[test]
+ fn from_vector_candidate_applies_distance_to_score() {
+ // L2: score = 1/(1+distance); distance 1.0 -> score 0.5 (score !=
distance).
+ let candidate = vector_candidate(1.0);
+ let position =
+ PrimaryKeySearchPosition::from_vector_candidate(&candidate,
VectorSearchMetric::L2)
+ .unwrap();
+ assert_eq!(position.score(), 0.5);
+ assert_eq!(position.row_position(), 0);
+ assert_eq!(position.data_file_name(), "f");
+ }
+
+ #[cfg(feature = "fulltext")]
+ #[test]
+ fn from_full_text_candidate_preserves_raw_score() {
+ use crate::table::pk_full_text_read::PrimaryKeyFullTextCandidate;
+ let candidate =
+ PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, 0.75,
"f".to_string(), 2)
+ .unwrap();
+ let position =
PrimaryKeySearchPosition::from_full_text_candidate(&candidate).unwrap();
+ assert_eq!(position.score(), 0.75);
+ assert_eq!(position.row_position(), 2);
+ }
+}
diff --git a/crates/paimon/src/table/pk_search_ranker.rs
b/crates/paimon/src/table/pk_search_ranker.rs
new file mode 100644
index 00000000..e3003a3d
--- /dev/null
+++ b/crates/paimon/src/table/pk_search_ranker.rs
@@ -0,0 +1,452 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Deterministic reciprocal-rank fusion for physical primary-key search
positions.
+//!
+//! Rust mirror of Java
+//! `org.apache.paimon.table.source.PrimaryKeySearchRanker`. It fuses the
+//! per-route rankings produced by the vector and full-text search routes
(each a
+//! list of [`PrimaryKeySearchPosition`]) into a single globally ranked list.
+//!
+//! Positions are combined ACROSS routes by their physical identity: the
+//! `PrimaryKeySearchPosition` [`Eq`]/[`Hash`] deliberately exclude the score,
so
+//! the same physical row hit by two routes collapses to one fused entry. A
+//! duplicate physical position WITHIN a single ranking is a caller error and
the
+//! weighted rankers fail loud, exactly like the Java `checkArgument`s.
+
+use std::cmp::Ordering;
+use std::collections::{HashMap, HashSet};
+
+use crate::table::pk_search_position::PrimaryKeySearchPosition;
+
+/// Reciprocal-rank-fusion smoothing constant. Matches Java's
+/// `PrimaryKeySearchRanker.DEFAULT_RRF_K` (an `int` with value `60`).
+pub(crate) const DEFAULT_RRF_K: i32 = 60;
+
+fn data_invalid(message: impl Into<String>) -> crate::Error {
+ crate::Error::DataInvalid {
+ message: message.into(),
+ source: None,
+ }
+}
+
+/// One locally scored ranking and its route weight. Mirrors Java
+/// `PrimaryKeySearchRanker.Ranking`.
+#[derive(Clone, Debug)]
+pub(crate) struct Ranking {
+ positions: Vec<PrimaryKeySearchPosition>,
+ weight: f64,
+}
+
+impl Ranking {
+ /// Build a ranking, rejecting a non-finite or non-positive `weight`.
+ /// Mirrors the `checkArgument` in the Java `Ranking` constructor.
+ pub(crate) fn new(
+ positions: Vec<PrimaryKeySearchPosition>,
+ weight: f64,
+ ) -> crate::Result<Self> {
+ if !weight.is_finite() || weight <= 0.0 {
+ return Err(data_invalid(format!(
+ "Search route weight must be finite and positive: {weight}."
+ )));
+ }
+ Ok(Self { positions, weight })
+ }
+
+ pub(crate) fn positions(&self) -> &[PrimaryKeySearchPosition] {
+ &self.positions
+ }
+
+ pub(crate) fn weight(&self) -> f64 {
+ self.weight
+ }
+}
+
+/// Java `PrimaryKeySearchPosition.compareTo`: partition bytes (unsigned,
+/// shorter-is-less on a common prefix) → bucket → data file name → row
position.
+fn compare_to(left: &PrimaryKeySearchPosition, right:
&PrimaryKeySearchPosition) -> Ordering {
+ left.partition()
+ .to_serialized_bytes()
+ .cmp(&right.partition().to_serialized_bytes())
+ .then_with(|| left.bucket().cmp(&right.bucket()))
+ .then_with(|| left.data_file_name().cmp(right.data_file_name()))
+ .then_with(|| left.row_position().cmp(&right.row_position()))
+}
+
+/// Java `LOCAL_BEST_FIRST`: score descending, then the physical tie-break.
+/// `total_cmp` reproduces `Float.compare` for the finite scores these
positions
+/// are validated to carry.
+fn best_first(left: &PrimaryKeySearchPosition, right:
&PrimaryKeySearchPosition) -> Ordering {
+ right
+ .score()
+ .total_cmp(&left.score())
+ .then_with(|| compare_to(left, right))
+}
+
+/// Selects globally highest-scored physical positions without rewriting their
+/// scores. Mirrors Java `topKByScore`.
+#[allow(dead_code)]
+pub(crate) fn top_k_by_score(
+ rankings: &[Vec<PrimaryKeySearchPosition>],
+ limit: usize,
+) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ check_limit(limit, "Search result limit")?;
+ let mut unique: HashMap<PrimaryKeySearchPosition,
PrimaryKeySearchPosition> = HashMap::new();
+ for ranking in rankings {
+ for position in ranking {
+ match unique.get(position) {
+ Some(previous) if best_first(position, previous) !=
Ordering::Less => {}
+ _ => {
+ unique.insert(position.clone(), position.clone());
+ }
+ }
+ }
+ }
+ let mut result: Vec<PrimaryKeySearchPosition> =
unique.into_values().collect();
+ result.sort_by(best_first);
+ result.truncate(limit);
+ Ok(result)
+}
+
+/// Convenience wrapper for equally weighted routes. Mirrors Java `rrf`.
+#[allow(dead_code)]
+pub(crate) fn rrf(
+ rankings: &[Vec<PrimaryKeySearchPosition>],
+ limit: usize,
+) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ let weighted: Vec<Ranking> = rankings
+ .iter()
+ .map(|ranking| Ranking::new(ranking.clone(), 1.0))
+ .collect::<crate::Result<_>>()?;
+ weighted_rrf(&weighted, limit)
+}
+
+/// Weighted reciprocal-rank fusion. Mirrors Java `weightedRrf`.
+pub(crate) fn weighted_rrf(
+ rankings: &[Ranking],
+ limit: usize,
+) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ check_limit(limit, "RRF result limit")?;
+ let mut fused_scores: HashMap<PrimaryKeySearchPosition, f64> =
HashMap::new();
+ for ranking in rankings {
+ add_ranking(&mut fused_scores, ranking)?;
+ }
+ top_k(fused_scores, limit)
+}
+
+/// Fuses heterogeneous route scores after independently normalizing each route
+/// to `[0, 1]`. Mirrors Java `weightedScore`.
+pub(crate) fn weighted_score(
+ rankings: &[Ranking],
+ limit: usize,
+) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ check_limit(limit, "Weighted-score result limit")?;
+ let mut fused_scores: HashMap<PrimaryKeySearchPosition, f64> =
HashMap::new();
+ for ranking in rankings {
+ let mut unique: HashSet<PrimaryKeySearchPosition> = HashSet::new();
+ let mut min = f32::INFINITY;
+ let mut max = f32::NEG_INFINITY;
+ for position in ranking.positions() {
+ if !unique.insert(position.clone()) {
+ return Err(data_invalid(format!(
+ "One weighted-score ranking contains duplicate physical
position {}.",
+ describe(position)
+ )));
+ }
+ min = min.min(position.score());
+ max = max.max(position.score());
+ }
+ let range = max - min;
+ for position in ranking.positions() {
+ let normalized: f64 = if range > 0.0 {
+ ((position.score() - min) / range) as f64
+ } else {
+ 1.0
+ };
+ *fused_scores.entry(position.clone()).or_insert(0.0) +=
ranking.weight() * normalized;
+ }
+ }
+ top_k(fused_scores, limit)
+}
+
+/// Fuses routes using weighted reciprocal rank without the RRF smoothing
+/// constant. Mirrors Java `weightedMrr`.
+pub(crate) fn weighted_mrr(
+ rankings: &[Ranking],
+ limit: usize,
+) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ check_limit(limit, "MRR result limit")?;
+ let mut fused_scores: HashMap<PrimaryKeySearchPosition, f64> =
HashMap::new();
+ for ranking in rankings {
+ let mut sorted = ranking.positions().to_vec();
+ sorted.sort_by(best_first);
+ let mut unique: HashSet<PrimaryKeySearchPosition> = HashSet::new();
+ for (i, position) in sorted.iter().enumerate() {
+ if !unique.insert(position.clone()) {
+ return Err(data_invalid(format!(
+ "One MRR ranking contains duplicate physical position {}.",
+ describe(position)
+ )));
+ }
+ *fused_scores.entry(position.clone()).or_insert(0.0) +=
+ ranking.weight() / (i as f64 + 1.0);
+ }
+ }
+ top_k(fused_scores, limit)
+}
+
+/// Java `checkArgument(limit > 0, ...)`.
+fn check_limit(limit: usize, what: &str) -> crate::Result<()> {
+ if limit == 0 {
+ return Err(data_invalid(format!("{what} must be positive: {limit}.")));
+ }
+ Ok(())
+}
+
+fn describe(position: &PrimaryKeySearchPosition) -> String {
+ format!(
+ "PrimaryKeySearchPosition{{bucket={}, dataFileName='{}',
rowPosition={}, score={}}}",
+ position.bucket(),
+ position.data_file_name(),
+ position.row_position(),
+ position.score()
+ )
+}
+
+/// Mirrors Java `addRanking`: sort a route best-first, assign 1-based ranks
that
+/// ties share, and accumulate `weight / (DEFAULT_RRF_K + rank)` per position.
+fn add_ranking(
+ fused_scores: &mut HashMap<PrimaryKeySearchPosition, f64>,
+ ranking: &Ranking,
+) -> crate::Result<()> {
+ let mut sorted = ranking.positions().to_vec();
+ sorted.sort_by(best_first);
+ let mut unique: HashSet<PrimaryKeySearchPosition> = HashSet::new();
+ let mut rank = 0i32;
+ let mut previous_score = f32::NAN;
+ for (i, position) in sorted.iter().enumerate() {
+ if !unique.insert(position.clone()) {
+ return Err(data_invalid(format!(
+ "One RRF ranking contains duplicate physical position {}.",
+ describe(position)
+ )));
+ }
+ if i == 0 || position.score().total_cmp(&previous_score) !=
Ordering::Equal {
+ rank = i as i32 + 1;
+ previous_score = position.score();
+ }
+ let contribution = ranking.weight() / f64::from(DEFAULT_RRF_K + rank);
+ *fused_scores.entry(position.clone()).or_insert(0.0) += contribution;
+ }
+ Ok(())
+}
+
+/// Mirrors Java `topK`: rewrite each fused key's score to its fused value and
+/// return the best `limit` positions, ordered best-first.
+fn top_k(
+ fused_scores: HashMap<PrimaryKeySearchPosition, f64>,
+ limit: usize,
+) -> crate::Result<Vec<PrimaryKeySearchPosition>> {
+ let mut result: Vec<PrimaryKeySearchPosition> = fused_scores
+ .into_iter()
+ .map(|(position, score)| position.with_score(score as f32))
+ .collect::<crate::Result<_>>()?;
+ result.sort_by(best_first);
+ result.truncate(limit);
+ Ok(result)
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::spec::BinaryRow;
+
+ /// Mirrors the Java test's `position(dataFileName, rowPosition, score)`
+ /// helper: `BinaryRow.EMPTY_ROW`, bucket `0`.
+ fn position(data_file_name: &str, row_position: i64, score: f32) ->
PrimaryKeySearchPosition {
+ PrimaryKeySearchPosition::new(
+ BinaryRow::new(0),
+ 0,
+ data_file_name.to_string(),
+ row_position,
+ score,
+ )
+ .unwrap()
+ }
+
+ fn ranking(positions: Vec<PrimaryKeySearchPosition>, weight: f64) ->
Ranking {
+ Ranking::new(positions, weight).unwrap()
+ }
+
+ fn file_names(positions: &[PrimaryKeySearchPosition]) -> Vec<&str> {
+ positions.iter().map(|p| p.data_file_name()).collect()
+ }
+
+ fn close(actual: f32, expected: f32) {
+ assert!(
+ (actual - expected).abs() < 0.000001,
+ "expected {expected}, got {actual}"
+ );
+ }
+
+ // Port of Java
`testFusesDuplicatePositionsAndAssignsTheSameRankToLocalTies`.
+ #[test]
+ fn fuses_duplicate_positions_and_assigns_the_same_rank_to_local_ties() {
+ let a = position("a", 0, 10.0);
+ let b = position("b", 0, 10.0);
+ let c = position("c", 0, 5.0);
+ let fused = rrf(
+ &[
+ vec![c.clone(), b.clone(), a.clone()],
+ vec![a.with_score(1.0).unwrap(), c.with_score(8.0).unwrap()],
+ ],
+ 3,
+ )
+ .unwrap();
+
+ assert_eq!(file_names(&fused), vec!["a", "c", "b"]);
+ close(fused[0].score(), (1.0 / 61.0 + 1.0 / 62.0) as f32);
+ close(fused[1].score(), (1.0 / 63.0 + 1.0 / 61.0) as f32);
+ close(fused[2].score(), (1.0 / 61.0) as f32);
+ }
+
+ // Port of Java `testUsesRouteWeightsAndDeterministicPhysicalTieBreaking`.
+ #[test]
+ fn uses_route_weights_and_deterministic_physical_tie_breaking() {
+ let a = position("a", 0, 1.0);
+ let b = position("b", 0, 1.0);
+ let fused = weighted_rrf(
+ &[ranking(vec![a.clone()], 2.0), ranking(vec![b.clone()], 2.0)],
+ 1,
+ )
+ .unwrap();
+
+ assert_eq!(fused.len(), 1);
+ assert_eq!(fused[0].data_file_name(), "a");
+ close(fused[0].score(), (2.0 / 61.0) as f32);
+ }
+
+ // Port of Java `testWeightedScoreNormalizesEachPhysicalRoute`.
+ #[test]
+ fn weighted_score_normalizes_each_physical_route() {
+ let a = position("a", 0, 0.0);
+ let b = position("b", 0, 10.0);
+ let c = position("c", 0, 0.0);
+
+ let fused = weighted_score(
+ &[
+ ranking(vec![a.clone(), b.clone()], 1.0),
+ ranking(vec![a.with_score(100.0).unwrap(), c.clone()], 2.0),
+ ],
+ 3,
+ )
+ .unwrap();
+
+ assert_eq!(file_names(&fused), vec!["a", "b", "c"]);
+ let scores: Vec<f32> = fused.iter().map(|p| p.score()).collect();
+ assert_eq!(scores, vec![2.0, 1.0, 0.0]);
+ }
+
+ #[test]
+ fn rrf_k_matches_java_value() {
+ assert_eq!(DEFAULT_RRF_K, 60);
+ }
+
+ // Negative and positive zero are distinct scores under `Float.compare`
+ // semantics, so they must NOT share a fused rank: +0.0 ranks ahead of
-0.0.
+ #[test]
+ fn signed_zero_scores_do_not_share_rank() {
+ let a = position("a", 0, 0.0); // +0.0, best-first
+ let b = position("b", 0, -0.0); // -0.0
+ let fused = weighted_rrf(&[ranking(vec![a, b], 1.0)], 2).unwrap();
+
+ assert_eq!(file_names(&fused), vec!["a", "b"]);
+ close(fused[0].score(), (1.0 / 61.0) as f32); // rank 1
+ close(fused[1].score(), (1.0 / 62.0) as f32); // rank 2 (would be 1/61
under `!=`)
+ }
+
+ #[test]
+ fn duplicate_within_ranking_fails_for_weighted_rankers() {
+ let a = position("a", 0, 3.0);
+ let dup = a.with_score(9.0).unwrap();
+
+ assert!(weighted_rrf(&[ranking(vec![a.clone(), dup.clone()], 1.0)],
3).is_err());
+ assert!(rrf(&[vec![a.clone(), dup.clone()]], 3).is_err());
+ assert!(weighted_score(&[ranking(vec![a.clone(), dup.clone()], 1.0)],
3).is_err());
+ assert!(weighted_mrr(&[ranking(vec![a, dup], 1.0)], 3).is_err());
+ }
+
+ #[test]
+ fn cross_ranking_combines_by_physical_key() {
+ // The same physical position in two routes fuses into one entry whose
+ // RRF score is the sum of both single-position contributions.
+ let a = position("a", 0, 5.0);
+ let fused = rrf(&[vec![a.clone()], vec![a.with_score(0.1).unwrap()]],
5).unwrap();
+ assert_eq!(fused.len(), 1);
+ assert_eq!(fused[0].data_file_name(), "a");
+ close(fused[0].score(), (1.0 / 61.0 + 1.0 / 61.0) as f32);
+ }
+
+ #[test]
+ fn weighted_mrr_uses_ordinal_rank() {
+ // Ranks are 1-based ordinals of the score-sorted list; scores 3,2,1
give
+ // MRR contributions 1/1, 1/2, 1/3 with weight 1.
+ let a = position("a", 0, 3.0);
+ let b = position("b", 0, 2.0);
+ let c = position("c", 0, 1.0);
+ let fused = weighted_mrr(&[ranking(vec![c, b, a], 1.0)], 3).unwrap();
+ assert_eq!(file_names(&fused), vec!["a", "b", "c"]);
+ close(fused[0].score(), 1.0);
+ close(fused[1].score(), 0.5);
+ close(fused[2].score(), (1.0 / 3.0) as f32);
+ }
+
+ #[test]
+ fn top_k_by_score_keeps_best_instance_without_rewriting() {
+ let a = position("a", 0, 2.0);
+ let b = position("b", 0, 9.0);
+ // Duplicate physical "a" across rankings: keep the higher raw score
(7).
+ let fused = top_k_by_score(
+ &[vec![a.clone(), b.clone()], vec![a.with_score(7.0).unwrap()]],
+ 5,
+ )
+ .unwrap();
+ assert_eq!(file_names(&fused), vec!["b", "a"]);
+ assert_eq!(fused[0].score(), 9.0);
+ assert_eq!(fused[1].score(), 7.0);
+ }
+
+ #[test]
+ fn non_positive_limit_fails() {
+ let a = position("a", 0, 1.0);
+ assert!(rrf(&[vec![a.clone()]], 0).is_err());
+ assert!(weighted_rrf(&[ranking(vec![a.clone()], 1.0)], 0).is_err());
+ assert!(weighted_score(&[ranking(vec![a.clone()], 1.0)], 0).is_err());
+ assert!(weighted_mrr(&[ranking(vec![a.clone()], 1.0)], 0).is_err());
+ assert!(top_k_by_score(&[vec![a]], 0).is_err());
+ }
+
+ #[test]
+ fn ranking_rejects_non_positive_or_non_finite_weight() {
+ let a = position("a", 0, 1.0);
+ assert!(Ranking::new(vec![a.clone()], 0.0).is_err());
+ assert!(Ranking::new(vec![a.clone()], -1.0).is_err());
+ assert!(Ranking::new(vec![a.clone()], f64::NAN).is_err());
+ assert!(Ranking::new(vec![a.clone()], f64::INFINITY).is_err());
+ assert!(Ranking::new(vec![a], 1.0).is_ok());
+ }
+}
diff --git a/crates/paimon/src/table/pk_vector_scan.rs
b/crates/paimon/src/table/pk_vector_scan.rs
index e6ba0c55..d9681d1c 100644
--- a/crates/paimon/src/table/pk_vector_scan.rs
+++ b/crates/paimon/src/table/pk_vector_scan.rs
@@ -190,6 +190,12 @@ impl BucketAccumulator {
/// The per-bucket search splits produced by planning.
pub(crate) struct PkVectorScanPlan {
+ // The snapshot the plan resolved during planning (pinned before the index
+ // manifest is read). It is authoritative even when planning yields zero
+ // searchable splits, so a cross-route consistency guard can require one
+ // pinned snapshot across routes. It is `0` only for a table with no
snapshot
+ // at all (never written), which also yields empty `splits`.
+ pub snapshot_id: i64,
pub splits: Vec<PkVectorSearchSplit>,
}
@@ -236,19 +242,28 @@ impl<'a> PkVectorScan<'a> {
if let Some(filter) = &self.filter {
read_builder.with_filter(filter.clone());
}
- let data_splits = read_builder
+ // Plan the data splits and capture the snapshot the scan pinned in one
+ // pass. The trace carries the resolved snapshot id even when the scan
+ // yields zero data splits, so the plan reports its real snapshot id
+ // (required by the cross-route snapshot-consistency guard) instead of
+ // deriving it from a first split that may not exist.
+ let (data_plan, trace) = read_builder
.new_scan()
.with_scan_all_files()
- .plan()
- .await?
- .splits()
- .to_vec();
-
- // No data files -> nothing to search.
- let Some(first_split) = data_splits.first() else {
- return Ok(PkVectorScanPlan { splits: Vec::new() });
+ .plan_with_trace()
+ .await?;
+ let data_splits = data_plan.splits().to_vec();
+
+ // No snapshot at all (table never written): nothing to search and no
+ // snapshot to pin. The empty split list makes every downstream
consumer
+ // treat this as "no candidates", so this is the only plan without a
real
+ // snapshot id.
+ let Some(snapshot_id) = trace.snapshot_id else {
+ return Ok(PkVectorScanPlan {
+ snapshot_id: 0,
+ splits: Vec::new(),
+ });
};
- let snapshot_id = first_split.snapshot_id();
let snapshot = snapshot_manager.get_snapshot(snapshot_id).await?;
// Index-manifest scan into filtered ANN payload tuples.
@@ -295,7 +310,10 @@ impl<'a> PkVectorScan<'a> {
}
let splits = plan_from_inputs(snapshot_id, data_splits, entries)?;
- Ok(PkVectorScanPlan { splits })
+ Ok(PkVectorScanPlan {
+ snapshot_id,
+ splits,
+ })
}
}
diff --git a/crates/paimon/src/table/vector_search_builder.rs
b/crates/paimon/src/table/vector_search_builder.rs
index 6f04065f..2c60aa9a 100644
--- a/crates/paimon/src/table/vector_search_builder.rs
+++ b/crates/paimon/src/table/vector_search_builder.rs
@@ -114,6 +114,24 @@ pub struct BatchVectorSearchBuilder<'a> {
filter: Option<Predicate>,
}
+/// The primary-key vector route's search output plus the source context a
later
+/// materialization (or a hybrid fusion across routes) needs. `candidates` are
the
+/// best-first hits; `splits` are the per-bucket source splits their
`split_index`
+/// refers into (the authority for re-associating a hit to its `DataFileMeta`);
+/// `snapshot_id` is the single snapshot the plan resolved during planning
+/// (authoritative even when the plan yields zero splits; `0` only for a table
+/// with no snapshot at all); `metric` is the resolved distance metric used to
+/// turn distances into scores. Produced by
+/// [`VectorSearchBuilder::search_pk_route`].
+pub(crate) struct PkVectorRouteResult {
+ pub(crate) candidates: Vec<PkVectorCandidate>,
+ // Read by the hybrid route consumer (fuses routes before materializing);
the
+ // materialized read reaches candidates/splits/metric directly.
+ pub(crate) snapshot_id: i64,
+ pub(crate) splits: Vec<PkVectorSearchSplit>,
+ pub(crate) metric: VectorSearchMetric,
+}
+
impl<'a> VectorSearchBuilder<'a> {
pub(crate) fn new(table: &'a Table) -> Self {
Self {
@@ -412,6 +430,39 @@ impl<'a> VectorSearchBuilder<'a> {
Ok((candidates.remove(0), plan, metric))
}
+ /// Plan + search the primary-key vector route and return the best-first
+ /// candidates together with the route source context needed to materialize
+ /// them later: the resolved snapshot id and the per-bucket source splits
+ /// (`split_index` is only meaningful against these originating splits),
plus
+ /// the resolved distance metric. This is the hybrid-reachable entry point
—
+ /// it runs exactly the plan/search core
[`execute_read`](Self::execute_read)
+ /// uses but stops before materialization, so a caller can fuse these
+ /// candidates before materializing. The materialized read is layered on
top
+ /// of it (see `execute_primary_key_vector_read`). The snapshot id is the
one
+ /// the plan pinned during planning — authoritative even for an empty plan
+ /// (which also yields empty candidates and empty splits); it is `0` only
for
+ /// a table with no snapshot at all.
+ pub(crate) async fn search_pk_route(
+ &self,
+ core: &CoreOptions<'_>,
+ pk_col: &str,
+ query_vector: &[f32],
+ limit: usize,
+ ) -> crate::Result<PkVectorRouteResult> {
+ let (candidates, plan, metric) = self
+ .plan_and_search_pk_candidates(core, pk_col, query_vector, limit)
+ .await?;
+ // Planning pins a single snapshot (`plan.snapshot_id`) even when it
yields
+ // zero searchable splits, so report it unconditionally rather than
deriving
+ // it from a split that may not exist.
+ Ok(PkVectorRouteResult {
+ candidates,
+ snapshot_id: plan.snapshot_id,
+ splits: plan.splits,
+ metric,
+ })
+ }
+
/// Materialize the best-first PK-vector search hits into Arrow rows.
Mirrors
/// Java `PrimaryKeyVectorRead` feeding its result splits into an ordinary
table
/// read: the search decides which rows, a subsequent read decides which
@@ -429,8 +480,13 @@ impl<'a> VectorSearchBuilder<'a> {
query_vector: &[f32],
limit: usize,
) -> crate::Result<ArrowRecordBatchStream> {
- let (candidates, plan, metric) = self
- .plan_and_search_pk_candidates(core, pk_col, query_vector, limit)
+ let PkVectorRouteResult {
+ candidates,
+ splits,
+ metric,
+ ..
+ } = self
+ .search_pk_route(core, pk_col, query_vector, limit)
.await?;
// Resolve the materialization read-type up front so an invalid
projection
@@ -451,7 +507,7 @@ impl<'a> VectorSearchBuilder<'a> {
Vec::new(),
);
- Self::materialize_candidates(candidates, &plan, metric,
&materialize_reader).await
+ Self::materialize_candidates(candidates, &splits, metric,
&materialize_reader).await
}
/// Materialize one best-first candidate list into an Arrow stream,
best-first,
@@ -462,7 +518,7 @@ impl<'a> VectorSearchBuilder<'a> {
/// use this so their materialization is identical.
async fn materialize_candidates(
candidates: Vec<PkVectorCandidate>,
- plan: &PkVectorScanPlan,
+ splits: &[PkVectorSearchSplit],
metric: VectorSearchMetric,
materialize_reader: &DataFileReader,
) -> crate::Result<ArrowRecordBatchStream> {
@@ -487,7 +543,7 @@ impl<'a> VectorSearchBuilder<'a> {
);
}
- let indexed_splits = build_indexed_splits(candidates, &plan.splits,
metric)?;
+ let indexed_splits = build_indexed_splits(candidates, splits, metric)?;
// Materialize every indexed split, retaining each batch and, per row,
the
// (rank, batch_index, row_index) tuple so we can reorder to
best-first.
@@ -564,12 +620,12 @@ fn is_reserved_read_column(name: &str) -> bool {
/// Reject a materialized read type whose resolved fields contain a reserved
/// metadata column name. Applied to the RESOLVED field list so the default
/// (all user columns) projection is covered, not only an explicit one.
-fn ensure_no_reserved_read_columns(fields: &[DataField]) -> crate::Result<()> {
+pub(crate) fn ensure_no_reserved_read_columns(fields: &[DataField]) ->
crate::Result<()> {
for field in fields {
if is_reserved_read_column(field.name()) {
return Err(crate::Error::DataInvalid {
message: format!(
- "vector search read projection must not include reserved
column '{}'",
+ "search read must not include reserved column '{}'",
field.name()
),
source: None,
@@ -1198,7 +1254,7 @@ impl<'a> BatchVectorSearchBuilder<'a> {
streams.push(
VectorSearchBuilder::materialize_candidates(
candidates,
- &plan,
+ &plan.splits,
metric,
&materialize_reader,
)
@@ -4420,6 +4476,257 @@ mod tests {
table
}
+ /// One Java `DataOutput#writeUTF` value (u16-BE length + modified UTF-8),
used
+ /// to assemble the `PrimaryKeyIndexSourceMeta` frame below.
+ fn java_write_utf(s: &str) -> Vec<u8> {
+ let mut body = Vec::new();
+ for c in s.encode_utf16() {
+ if (0x0001..=0x007F).contains(&c) {
+ body.push(c as u8);
+ } else if c > 0x07FF {
+ body.push(0xE0 | (c >> 12) as u8);
+ body.push(0x80 | ((c >> 6) & 0x3F) as u8);
+ body.push(0x80 | (c & 0x3F) as u8);
+ } else {
+ body.push(0xC0 | (c >> 6) as u8);
+ body.push(0x80 | (c & 0x3F) as u8);
+ }
+ }
+ let mut out = (body.len() as u16).to_be_bytes().to_vec();
+ out.extend_from_slice(&body);
+ out
+ }
+
+ /// The Java `PrimaryKeyIndexSourceMeta` frame: `i32-BE version=1`, `i32-BE
+ /// data_level`, `i32-BE count`, then per source file a `writeUTF` name
and an
+ /// `i64-BE` row count.
+ fn pk_source_meta_bytes(data_level: i32, files: &[(&str, i64)]) -> Vec<u8>
{
+ let mut out = Vec::new();
+ out.extend_from_slice(&1i32.to_be_bytes());
+ out.extend_from_slice(&data_level.to_be_bytes());
+ out.extend_from_slice(&(files.len() as i32).to_be_bytes());
+ for (name, rows) in files {
+ out.extend_from_slice(&java_write_utf(name));
+ out.extend_from_slice(&rows.to_be_bytes());
+ }
+ out
+ }
+
+ /// Build a committed primary-key vector table (memory FS) over `vectors`
+ /// (dimension 2): write a real data file via the write path, promote its
meta
+ /// to a compacted, non-level-0 file (the PK index-source precondition),
then
+ /// build + commit a real vindex IVF-flat ANN segment naming that file.
Single
+ /// bucket, `nlist = 1`, so the ANN search is exact. Returns the opened
table,
+ /// ready for `search_pk_route`.
+ async fn build_committed_pk_vector_table(vectors: &[[f32; 2]]) -> Table {
+ use crate::spec::{GlobalIndexMeta, IndexFileMeta, VectorType};
+ use crate::table::CommitMessage;
+ use bytes::Bytes;
+ use paimon_vindex_core::index::{VectorIndexConfig, VectorIndexTrainer,
VectorIndexWriter};
+ use paimon_vindex_core::io::PosWriter;
+
+ const DIM: usize = 2;
+ let table_path = "memory:/pk_vector_route_test";
+ let schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(
+ "embedding",
+ DataType::Vector(
+ VectorType::try_new(true, DIM as u32,
DataType::Float(FloatType::new()))
+ .unwrap(),
+ ),
+ )
+ .primary_key(["id"])
+ .option("bucket", "1")
+ .option("pk-vector.index.columns", "embedding")
+ .option("fields.embedding.pk-vector.index.type",
IVF_FLAT_IDENTIFIER)
+ .option("fields.embedding.pk-vector.distance.metric", "l2")
+ .build()
+ .unwrap();
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let table = Table::new(
+ file_io.clone(),
+ Identifier::new("default", "pk_vector_route_test"),
+ table_path.to_string(),
+ TableSchema::new(0, &schema),
+ None,
+ );
+ for dir in ["snapshot", "manifest", "index"] {
+ file_io
+ .mkdirs(&format!("{table_path}/{dir}"))
+ .await
+ .unwrap();
+ }
+
+ // id + FixedSizeList<Float32> batch matching the table's target
schema.
+ let ids: Vec<i32> = (0..vectors.len() as i32).collect();
+ let element_field = Arc::new(ArrowField::new("element",
ArrowDataType::Float32, true));
+ let mut vec_builder = FixedSizeListBuilder::new(Float32Builder::new(),
DIM as i32)
+ .with_field(element_field.clone());
+ for v in vectors {
+ for &x in v {
+ vec_builder.values().append_value(x);
+ }
+ vec_builder.append(true);
+ }
+ let arrow_schema = Arc::new(ArrowSchema::new(vec![
+ ArrowField::new("id", ArrowDataType::Int32, false),
+ ArrowField::new(
+ "embedding",
+ ArrowDataType::FixedSizeList(element_field, DIM as i32),
+ true,
+ ),
+ ]));
+ let batch = RecordBatch::try_new(
+ arrow_schema,
+ vec![
+ Arc::new(Int32Array::from(ids)) as ArrayRef,
+ Arc::new(vec_builder.finish()) as ArrayRef,
+ ],
+ )
+ .unwrap();
+
+ // Real data-file meta via the write path (these messages are not
committed
+ // as-is; the meta is promoted below and committed with the index).
+ let mut writer = TableWrite::new(&table,
"route-test".to_string()).unwrap();
+ writer.write_arrow_batch(&batch).await.unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ let base = &messages[0];
+ let base_meta = base.new_files[0].clone();
+ let bucket = base.bucket;
+ let partition = base.partition.clone();
+ let data_file_name = base_meta.file_name.clone();
+ let row_count = base_meta.row_count;
+
+ // PK index-source precondition: compacted, non-level-0, first_row_id
pinned.
+ let indexed_meta = DataFileMeta {
+ level: 1,
+ file_source: Some(1),
+ first_row_id: Some(0),
+ ..base_meta
+ };
+
+ // Real vindex IVF-flat segment (nlist=1 -> exact) over the vectors.
+ let n = vectors.len();
+ let flat: Vec<f32> = vectors.iter().flat_map(|v|
v.iter().copied()).collect();
+ let seg_ids: Vec<i64> = (0..n as i64).collect();
+ let native = HashMap::from([
+ ("index.type".to_string(), "ivf_flat".to_string()),
+ ("dimension".to_string(), DIM.to_string()),
+ ("nlist".to_string(), "1".to_string()),
+ ("metric".to_string(), "l2".to_string()),
+ ]);
+ let config = VectorIndexConfig::from_options(&native).unwrap();
+ let training = VectorIndexTrainer::train(config, &flat, n).unwrap();
+ let mut ann_writer = VectorIndexWriter::new(training);
+ ann_writer.add_vectors(&seg_ids, &flat, n).unwrap();
+ let mut seg_bytes = Vec::new();
+ {
+ let mut out = PosWriter::new(&mut seg_bytes);
+ ann_writer.write(&mut out).unwrap();
+ }
+ let index_file_name = "vector-ivf-flat-route.index".to_string();
+ let index_file_size = seg_bytes.len() as u64;
+ file_io
+ .new_output(&format!("{table_path}/index/{index_file_name}"))
+ .unwrap()
+ .write(Bytes::from(seg_bytes))
+ .await
+ .unwrap();
+
+ let vector_field_id = schema
+ .fields()
+ .iter()
+ .find(|f| f.name() == "embedding")
+ .unwrap()
+ .id();
+ let index_file = IndexFileMeta {
+ index_type: IVF_FLAT_IDENTIFIER.to_string(),
+ file_name: index_file_name,
+ file_size: i64::try_from(index_file_size).unwrap(),
+ row_count: i32::try_from(row_count).unwrap(),
+ deletion_vectors_ranges: None,
+ global_index_meta: Some(GlobalIndexMeta {
+ row_range_start: 0,
+ row_range_end: row_count - 1,
+ index_field_id: vector_field_id,
+ extra_field_ids: None,
+ source_meta: Some(pk_source_meta_bytes(1, &[(&data_file_name,
row_count)])),
+ index_meta: None,
+ }),
+ };
+
+ let mut message = CommitMessage::new(partition, bucket,
vec![indexed_meta]);
+ message.new_index_files = vec![index_file];
+ TableCommit::new(table.clone(), "route-test".to_string())
+ .commit(vec![message])
+ .await
+ .unwrap();
+ table
+ }
+
+ // ---- search_pk_route: candidate-only producer returns candidates +
context ----
+ #[tokio::test]
+ async fn search_pk_route_returns_candidates_and_source_context() {
+ // query [0,1]: squared-L2 distances pos1=0 < pos2=1 < pos0=2, so the
+ // strict-gap top-2 is [pos1, pos2] (best-first, not physical order).
+ let table = build_committed_pk_vector_table(&[[1.0, 0.0], [0.0, 1.0],
[1.0, 1.0]]).await;
+ let core = CoreOptions::new(table.schema().options());
+ let builder = table.new_vector_search_builder();
+ let route = builder
+ .search_pk_route(&core, "embedding", &[0.0, 1.0], 2)
+ .await
+ .unwrap();
+
+ // Two nearest neighbours recalled, best-first, without
materialization.
+ assert_eq!(route.candidates.len(), 2, "top-2 candidates expected");
+ assert_eq!(route.candidates[0].row_position, 1, "nearest is position
1");
+ assert_eq!(
+ route.candidates[1].row_position, 2,
+ "second nearest is position 2"
+ );
+ assert!(
+ route.candidates[0].distance <= route.candidates[1].distance,
+ "candidates must be best-first by distance"
+ );
+
+ // Source context present for a non-empty plan: the snapshot the plan
+ // pinned during planning (a real id, not None/0), per-bucket source
+ // splits, and the resolved metric.
+ assert_eq!(route.snapshot_id, 1, "first commit -> snapshot 1");
+ assert!(!route.splits.is_empty(), "non-empty source splits expected");
+ assert_eq!(route.metric, VectorSearchMetric::L2);
+ assert!(
+ route
+ .candidates
+ .iter()
+ .all(|c| c.split_index < route.splits.len()),
+ "candidate split_index must refer into the returned splits"
+ );
+ }
+
+ /// A table with no snapshot at all (never written) yields empty candidates
+ /// and empty source splits; with no snapshot to pin the id is `0`, and the
+ /// metric still resolves.
+ #[tokio::test]
+ async fn search_pk_route_empty_plan_yields_empty_context() {
+ 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 core = CoreOptions::new(table.schema().options());
+ let builder = table.new_vector_search_builder();
+ let route = builder
+ .search_pk_route(&core, "embedding", &[1.0f32; 128], 3)
+ .await
+ .unwrap();
+ assert!(route.candidates.is_empty(), "no data -> no candidates");
+ assert!(route.splits.is_empty(), "no data -> no source splits");
+ assert_eq!(route.snapshot_id, 0, "no snapshot -> zero id");
+ assert_eq!(route.metric, VectorSearchMetric::L2);
+ }
+
#[tokio::test]
async fn pk_branch_disabled_falls_through_to_de_path() {
// No pk-vector.index.columns: behaves exactly as the DE path. With no