From 2448ae31770c8bf13c268e428e92e2337eeb8571 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Thu, 6 Aug 2026 15:14:47 +0000 Subject: [PATCH 1/7] fix(index): filter stale segment rows before vector top-k --- python/python/tests/test_vector_index.py | 64 ++++++ rust/lance/src/index.rs | 4 +- rust/lance/src/io/exec/knn.rs | 248 ++++++++++++++++++++++- 3 files changed, 305 insertions(+), 11 deletions(-) diff --git a/python/python/tests/test_vector_index.py b/python/python/tests/test_vector_index.py index 21ae2aaac8e..5b95e898363 100644 --- a/python/python/tests/test_vector_index.py +++ b/python/python/tests/test_vector_index.py @@ -1961,6 +1961,70 @@ def test_optimize_indices(indexed_dataset): assert stats["num_indices"] == 2 +@pytest.mark.parametrize("enable_stable_row_ids", [False, True]) +def test_segment_ownership_filter_precedes_partition_topk( + tmp_path, enable_stable_row_ids +): + ndim = 4 + + def table(ids, value): + vectors = np.full((len(ids), ndim), value, dtype=np.float32) + return pa.table( + { + "id": pa.array(ids, type=pa.int64()), + "vector": pa.FixedSizeListArray.from_arrays( + pa.array(vectors.reshape(-1), type=pa.float32()), ndim + ), + } + ) + + dataset = lance.write_dataset( + table(range(20), 1.0), + tmp_path, + mode="create", + enable_stable_row_ids=enable_stable_row_ids, + ) + dataset = lance.write_dataset( + table(range(100, 120), 0.0), dataset.uri, mode="append" + ) + dataset = dataset.create_index( + "vector", index_type="IVF_FLAT", metric="l2", num_partitions=1 + ) + + fragment = dataset.get_fragment(1) + row_ids = fragment.to_table(columns=["id"], with_row_id=True)["_rowid"].to_pylist() + update_data = pa.table( + { + "_rowid": pa.array(row_ids, type=pa.uint64()), + "vector": pa.array( + [[10.0] * ndim] * len(row_ids), type=pa.list_(pa.float32(), ndim) + ), + } + ) + updated_fragment, fields_modified = fragment.update_columns(update_data) + dataset = lance.LanceDataset.commit( + dataset.uri, + lance.LanceOperation.Update( + updated_fragments=[updated_fragment], fields_modified=fields_modified + ), + read_version=dataset.version, + ) + dataset.optimize.optimize_indices(num_indices_to_merge=0) + dataset = lance.dataset(dataset.uri) + + result = dataset.to_table( + columns=["id"], + nearest={ + "column": "vector", + "q": np.zeros(ndim, dtype=np.float32), + "k": 5, + }, + ) + + assert all(row_id < 20 for row_id in result["id"].to_pylist()) + assert result["_distance"].to_pylist() == pytest.approx([4.0] * 5) + + @pytest.mark.skip(reason="retrain is deprecated") def test_retrain_indices(indexed_dataset): data = create_table() diff --git a/rust/lance/src/index.rs b/rust/lance/src/index.rs index adea488689b..887ce02acb8 100644 --- a/rust/lance/src/index.rs +++ b/rust/lance/src/index.rs @@ -132,14 +132,14 @@ fn validate_segment_metadata(index_name: &str, segments: &[IndexMetadata]) -> Re Ok(()) } -fn collect_subtree_field_ids(field: &Field, field_ids: &mut HashSet) { +pub(crate) fn collect_subtree_field_ids(field: &Field, field_ids: &mut HashSet) { field_ids.insert(field.id); for child in &field.children { collect_subtree_field_ids(child, field_ids); } } -fn fragment_field_paths<'a>( +pub(crate) fn fragment_field_paths<'a>( fragment: &'a Fragment, indexed_field_ids: &HashSet, ) -> HashMap { diff --git a/rust/lance/src/io/exec/knn.rs b/rust/lance/src/io/exec/knn.rs index e5157f78818..f952d5508f8 100644 --- a/rust/lance/src/io/exec/knn.rs +++ b/rust/lance/src/io/exec/knn.rs @@ -63,9 +63,9 @@ use uuid::Uuid; use lance_select::RowAddrMask; use crate::dataset::Dataset; -use crate::index::DatasetIndexInternalExt; use crate::index::prefilter::{DatasetPreFilter, FilterLoader}; use crate::index::vector::utils::{get_vector_type, validate_distance_type_for}; +use crate::index::{DatasetIndexInternalExt, collect_subtree_field_ids, fragment_field_paths}; use crate::{Error, Result}; use lance_arrow::*; @@ -1545,6 +1545,133 @@ struct LatePartitionSearchControl { max_results: usize, } +/// A query prefilter restricted to the fragments owned by one physical index segment. +/// +/// The shared dataset prefilter covers the union of every segment. When an in-place +/// update removes a fragment from an older segment's metadata, this additional mask +/// keeps that segment's stale physical rows out of its local top-k. +struct SegmentPreFilter { + base: Arc, + ownership_mask: Arc, + final_mask: Mutex>>, +} + +impl SegmentPreFilter { + fn new(base: Arc, ownership_mask: Arc) -> Self { + Self { + base, + ownership_mask, + final_mask: Mutex::new(None), + } + } +} + +#[async_trait::async_trait] +impl PreFilter for SegmentPreFilter { + async fn wait_for_ready(&self) -> Result<()> { + self.base.wait_for_ready().await?; + let mut final_mask = self.final_mask.lock().unwrap(); + final_mask.get_or_insert_with(|| { + Arc::new(self.base.mask().as_ref().clone() & self.ownership_mask.as_ref().clone()) + }); + Ok(()) + } + + fn is_empty(&self) -> bool { + false + } + + fn mask(&self) -> Arc { + self.final_mask + .lock() + .unwrap() + .as_ref() + .expect("mask called without call to wait_for_ready") + .clone() + } + + fn filter_row_ids<'a>(&self, row_ids: Box + 'a>) -> Vec { + self.mask().selected_indices(row_ids) + } +} + +/// Return whether a segment can still physically contain rows from a fragment it no +/// longer owns. Append-only deltas have disjoint bitmaps too, so the bitmap alone is +/// not enough: compare indexed data files with the snapshot used to build the segment. +async fn segment_needs_ownership_filter(dataset: &Dataset, index: &IndexMetadata) -> bool { + let Some(owned_fragments) = index.fragment_bitmap.as_ref() else { + return false; + }; + if dataset + .fragment_bitmap + .iter() + .all(|fragment_id| owned_fragments.contains(fragment_id)) + { + return false; + } + + let historical = match dataset.checkout_version(index.dataset_version).await { + Ok(historical) => historical, + Err(error) => { + // Old manifests may have been cleaned up. Restrict conservatively so a + // missing historical snapshot cannot make stale index rows visible again. + log::debug!( + "Could not inspect dataset version {} for index segment {}: {}. Applying the segment ownership filter conservatively.", + index.dataset_version, + index.uuid, + error + ); + return true; + } + }; + + let mut indexed_field_ids = HashSet::new(); + for field_id in &index.fields { + if let Some(field) = dataset.schema().field_by_id(*field_id) { + collect_subtree_field_ids(field, &mut indexed_field_ids); + } + } + let current_fragments = dataset + .fragments() + .iter() + .map(|fragment| (fragment.id as u32, fragment)) + .collect::>(); + + historical.fragments().iter().any(|historical_fragment| { + let fragment_id = historical_fragment.id as u32; + if owned_fragments.contains(fragment_id) { + return false; + } + current_fragments + .get(&fragment_id) + .is_some_and(|current_fragment| { + fragment_field_paths(historical_fragment, &indexed_field_ids) + != fragment_field_paths(current_fragment, &indexed_field_ids) + }) + }) +} + +async fn prefilter_for_segment( + dataset: Arc, + index: &IndexMetadata, + base: Arc, +) -> Result> { + if !segment_needs_ownership_filter(dataset.as_ref(), index).await { + return Ok(base); + } + + let Some(owned_fragments) = index.fragment_bitmap.clone() else { + return Ok(base); + }; + let Some(ownership_mask) = + DatasetPreFilter::create_restricted_deletion_mask(dataset, owned_fragments) + else { + return Ok(base); + }; + let ownership_mask = ownership_mask.await?; + Ok(Arc::new(SegmentPreFilter::new(base, ownership_mask))) +} + impl PartitionSearchControl for LatePartitionSearchControl { fn should_stop(&self) -> bool { self.state.num_results_found.load(Ordering::Relaxed) >= self.max_results @@ -1589,7 +1716,7 @@ impl ANNIvfSubIndexExec { index: Arc, query: Query, part_id: usize, - pre_filter: Arc, + pre_filter: Arc, metrics: Arc, ) -> DataFusionResult { let batch = index @@ -1630,7 +1757,8 @@ impl ANNIvfSubIndexExec { query: Query, partitions: Arc, q_c_dists: Arc, - prefilter: Arc, + prefilter: Arc, + global_prefilter: Arc, metrics: Arc, state: Arc, target_partitions: usize, @@ -1654,7 +1782,7 @@ impl ANNIvfSubIndexExec { // We know the prefilter should be ready at this point so we shouldn't // need to call wait_for_ready - let prefilter_mask = prefilter.mask(); + let prefilter_mask = global_prefilter.mask(); let max_results = prefilter_mask.max_len().map(|x| x as usize); @@ -1781,7 +1909,7 @@ impl ANNIvfSubIndexExec { query: Query, partitions: Arc, q_c_dists: Arc, - prefilter: Arc, + prefilter: Arc, metrics: Arc, state: Arc, target_partitions: usize, @@ -1987,6 +2115,13 @@ impl ExecutionPlan for ANNIvfSubIndexExec { } Arc::new(pf) }; + let indices_by_uuid = Arc::new( + indices + .iter() + .cloned() + .map(|index| (index.uuid, index)) + .collect::>(), + ); let state = Arc::new(ANNIvfEarlySearchResults::new(indices.len(), query.k)); @@ -1998,11 +2133,23 @@ impl ExecutionPlan for ANNIvfSubIndexExec { let column = column.clone(); let metrics = metrics.clone(); let pre_filter = pre_filter.clone(); + let indices_by_uuid = indices_by_uuid.clone(); let state = state.clone(); let mut query = query.clone(); let pruned_nprobes = early_pruning(q_c_dists.values(), query.k); adjust_probes(&mut query, pruned_nprobes); async move { + let index_metadata = indices_by_uuid.get(&index_uuid).ok_or_else(|| { + DataFusionError::Execution(format!( + "ANNSubIndexExec: input referenced unknown index segment {index_uuid}" + )) + })?; + let segment_pre_filter = prefilter_for_segment( + ds.clone(), + index_metadata, + pre_filter.clone(), + ) + .await?; let raw_index = ds .open_vector_index(&column, &index_uuid, &metrics.index_metrics) .await?; @@ -2013,7 +2160,7 @@ impl ExecutionPlan for ANNIvfSubIndexExec { query.clone(), part_ids.clone(), q_c_dists.clone(), - pre_filter.clone(), + segment_pre_filter.clone(), metrics.clone(), state.clone(), target_partitions, @@ -2023,6 +2170,7 @@ impl ExecutionPlan for ANNIvfSubIndexExec { query, part_ids, q_c_dists, + segment_pre_filter, pre_filter, metrics, state, @@ -2301,8 +2449,8 @@ mod tests { use arrow::compute::{concat_batches, sort_to_indices, take_record_batch}; use arrow::datatypes::Float32Type; use arrow_array::{ - ArrayRef, FixedSizeListArray, Float32Array, Int32Array, RecordBatchIterator, StringArray, - StructArray, + ArrayRef, FixedSizeListArray, Float32Array, Int32Array, RecordBatchIterator, + RecordBatchReader, StringArray, StructArray, }; use arrow_schema::{Field as ArrowField, Schema as ArrowSchema}; use async_trait::async_trait; @@ -2831,6 +2979,86 @@ mod tests { prefilter } + #[tokio::test] + async fn test_append_only_deltas_keep_empty_prefilter_fast_path() { + let first = lance_datagen::gen_batch() + .col( + "vector", + array::rand_vec::(lance_datagen::Dimension::from(4)), + ) + .into_reader_rows(RowCount::from(20), BatchCount::from(1)); + let first_schema = first.schema(); + let mut dataset = Dataset::write(first, "memory://", None).await.unwrap(); + let first_version = dataset.manifest.version; + let first_fragments = dataset.fragment_bitmap.as_ref().clone(); + let field_id = dataset.schema().field("vector").unwrap().id; + + let second = lance_datagen::gen_batch() + .col( + "vector", + array::rand_vec::(lance_datagen::Dimension::from(4)), + ) + .into_reader_rows(RowCount::from(20), BatchCount::from(1)); + assert_eq!(second.schema(), first_schema); + dataset = Dataset::write( + second, + "memory://", + Some(WriteParams { + mode: WriteMode::Append, + ..Default::default() + }), + ) + .await + .unwrap(); + let appended_fragments = dataset.fragment_bitmap.as_ref() - &first_fragments; + let dataset = Arc::new(dataset); + + let old_segment = IndexMetadata { + uuid: Uuid::new_v4(), + fields: vec![field_id], + name: "vector_idx".to_string(), + dataset_version: first_version, + fragment_bitmap: Some(first_fragments), + index_details: None, + index_version: 0, + created_at: None, + base_id: None, + files: None, + }; + let new_segment = IndexMetadata { + uuid: Uuid::new_v4(), + fields: vec![field_id], + name: "vector_idx".to_string(), + dataset_version: dataset.manifest.version, + fragment_bitmap: Some(appended_fragments), + index_details: None, + index_version: 0, + created_at: None, + base_id: None, + files: None, + }; + assert!( + !segment_needs_ownership_filter(dataset.as_ref(), &old_segment).await, + "append-only data must not look like an in-place fragment update" + ); + let base = Arc::new(DatasetPreFilter::new( + dataset.clone(), + &[old_segment.clone(), new_segment], + None, + )); + base.wait_for_ready().await.unwrap(); + assert!(base.is_empty(), "the combined delta coverage is unfiltered"); + let segment_prefilter = prefilter_for_segment(dataset, &old_segment, base) + .await + .unwrap(); + segment_prefilter.wait_for_ready().await.unwrap(); + + assert!( + segment_prefilter.is_empty(), + "an append-only delta must preserve the unfiltered sub-index fast path" + ); + } + fn prepared_metrics() -> Arc { Arc::new(AnnIndexMetrics::new(&ExecutionPlanMetricsSet::new(), 0)) } @@ -3024,12 +3252,14 @@ mod tests { .unwrap(), ); + let prefilter = empty_prefilter().await; let batches = ANNIvfSubIndexExec::late_search( index, query, Arc::new(UInt32Array::from(vec![0, 1, 2])), Arc::new(Float32Array::from(vec![0.1, 0.2, 0.3])), - empty_prefilter().await, + prefilter.clone(), + prefilter, prepared_metrics(), state.clone(), usize::MAX, From 556894e7f6ac6f2bf3e105413d450e42c891ab2a Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Thu, 6 Aug 2026 16:23:38 +0000 Subject: [PATCH 2/7] fix(index): preserve vector segment physical coverage --- protos/index.proto | 8 +- rust/lance/src/dataset/optimize/remapping.rs | 14 ++- rust/lance/src/index.rs | 18 ++- rust/lance/src/index/append.rs | 5 + rust/lance/src/index/create.rs | 27 ++-- rust/lance/src/index/vector.rs | 14 ++- rust/lance/src/index/vector/details.rs | 124 +++++++++++++++++-- rust/lance/src/index/vector/ivf.rs | 7 +- rust/lance/src/index/vector/ivf/v2.rs | 10 ++ rust/lance/src/io/exec/knn.rs | 113 ++++++++--------- 10 files changed, 253 insertions(+), 87 deletions(-) diff --git a/protos/index.proto b/protos/index.proto index a72207b59fe..80af6db5130 100644 --- a/protos/index.proto +++ b/protos/index.proto @@ -225,6 +225,12 @@ message VectorIndexDetails { // Keys use reverse-DNS namespacing (e.g., "lance.ivf.max_iters", "lancedb.accelerator"). // Unrecognized keys must be silently ignored by all runtimes. map runtime_hints = 9; + + // The fragments whose rows are physically present in this segment. This is + // immutable even when fragment ownership is later pruned after an in-place + // update. It contains a serialized 32-bit Roaring bitmap when known and is + // absent for segments written before physical coverage was recorded. + optional bytes physical_fragment_bitmap = 10; } // Hierarchical Navigable Small World (HNSW) parameters, used as an optional configuration for IVF indexes. @@ -248,4 +254,4 @@ message BloomFilterIndexDetails {} message RTreeIndexDetails {} -message FMIndexDetails {} \ No newline at end of file +message FMIndexDetails {} diff --git a/rust/lance/src/dataset/optimize/remapping.rs b/rust/lance/src/dataset/optimize/remapping.rs index 9a021a44bb0..6bc5c473b29 100644 --- a/rust/lance/src/dataset/optimize/remapping.rs +++ b/rust/lance/src/dataset/optimize/remapping.rs @@ -8,6 +8,7 @@ use crate::Result; use crate::dataset::transaction::{Operation, Transaction}; use crate::index::DatasetIndexExt; use crate::index::frag_reuse::{load_frag_reuse_index_details, open_frag_reuse_index}; +use crate::index::vector::details::with_physical_fragment_bitmap; use crate::{Dataset, index}; use async_trait::async_trait; use lance_core::Error; @@ -356,7 +357,13 @@ async fn remap_index(dataset: &mut Dataset, index_id: &Uuid) -> Result<()> { fields: curr_index_meta.fields.clone(), dataset_version: new_dataset_version, fragment_bitmap: bitmap_after_remap, - index_details: curr_index_meta.index_details.clone(), + index_details: curr_index_meta + .index_details + .as_deref() + .cloned() + .map(|details| with_physical_fragment_bitmap(details, None)) + .transpose()? + .map(Arc::new), index_version: curr_index_meta.index_version, created_at: curr_index_meta.created_at, base_id: None, @@ -368,7 +375,10 @@ async fn remap_index(dataset: &mut Dataset, index_id: &Uuid) -> Result<()> { fields: curr_index_meta.fields.clone(), dataset_version: new_dataset_version, fragment_bitmap: bitmap_after_remap, - index_details: Some(Arc::new(remapped_index.index_details)), + index_details: Some(Arc::new(with_physical_fragment_bitmap( + remapped_index.index_details, + None, + )?)), index_version: remapped_index.index_version as i32, created_at: curr_index_meta.created_at, base_id: None, diff --git a/rust/lance/src/index.rs b/rust/lance/src/index.rs index 887ce02acb8..163ab3f1c40 100644 --- a/rust/lance/src/index.rs +++ b/rust/lance/src/index.rs @@ -64,8 +64,8 @@ use serde_json::json; use tracing::{info, instrument, warn}; use uuid::Uuid; use vector::details::{ - derive_vector_index_type, infer_missing_vector_details, needs_vector_details_inference, - vector_details_as_json, + derive_vector_index_type, infer_missing_vector_details, merged_physical_fragment_bitmap, + needs_vector_details_inference, vector_details_as_json, with_physical_fragment_bitmap, }; pub(crate) use vector::details::{vector_index_details, vector_index_details_default}; use vector::ivf::v2::{IVFIndex, IvfStateEntryBox}; @@ -132,14 +132,14 @@ fn validate_segment_metadata(index_name: &str, segments: &[IndexMetadata]) -> Re Ok(()) } -pub(crate) fn collect_subtree_field_ids(field: &Field, field_ids: &mut HashSet) { +fn collect_subtree_field_ids(field: &Field, field_ids: &mut HashSet) { field_ids.insert(field.id); for child in &field.children { collect_subtree_field_ids(child, field_ids); } } -pub(crate) fn fragment_field_paths<'a>( +fn fragment_field_paths<'a>( fragment: &'a Fragment, indexed_field_ids: &HashSet, ) -> HashMap { @@ -1947,13 +1947,21 @@ impl DatasetIndexExt for Dataset { }; let last_idx = deltas.last().expect("Delta indices should not be empty"); + let physical_fragment_bitmap = merged_physical_fragment_bitmap( + res.removed_indices.iter().copied(), + &res.new_fragment_bitmap, + ); + let new_index_details = with_physical_fragment_bitmap( + res.new_index_details, + physical_fragment_bitmap.as_ref(), + )?; let new_idx = IndexMetadata { uuid: res.new_uuid, name: last_idx.name.clone(), // Keep the same name fields: last_idx.fields.clone(), dataset_version: res.new_dataset_version, fragment_bitmap: Some(res.new_fragment_bitmap), - index_details: Some(Arc::new(res.new_index_details)), + index_details: Some(Arc::new(new_index_details)), index_version: res.new_index_version, created_at: Some(chrono::Utc::now()), base_id: None, // New merged index file locates in the cloned dataset. diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 1e36f4b1f75..398d8ee8170 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -2586,6 +2586,11 @@ mod tests { dataset.fragment_bitmap.as_ref(), "the compatible merge must preserve exact fragment coverage" ); + assert_eq!( + crate::index::vector::details::physical_fragment_bitmap(&merged[0]), + Some(dataset.fragment_bitmap.as_ref().clone()), + "the compatible merge must preserve physical fragment coverage" + ); } #[tokio::test] diff --git a/rust/lance/src/index/create.rs b/rust/lance/src/index/create.rs index 16bb655a871..e954cb88d11 100644 --- a/rust/lance/src/index/create.rs +++ b/rust/lance/src/index/create.rs @@ -14,6 +14,7 @@ use crate::{ vector::{ LANCE_VECTOR_INDEX, StageParams, VectorIndexParams, build_distributed_vector_index, build_empty_vector_index, build_filtered_vector_index, build_vector_index, + details::with_physical_fragment_bitmap, }, vector_index_details, vector_index_details_default, }, @@ -592,21 +593,27 @@ impl<'a> CreateIndexBuilder<'a> { } }; + let fragment_bitmap = if train { + match &self.fragments { + Some(fragment_ids) => fragment_ids.iter().collect(), + None => self.dataset.fragment_bitmap.as_ref().clone(), + } + } else { + // Empty bitmap for untrained indices + roaring::RoaringBitmap::new() + }; + let index_details = with_physical_fragment_bitmap( + created_index.index_details, + Some(&fragment_bitmap), + )?; + Ok(IndexMetadata { uuid: output_index_uuid, name: index_name, fields: vec![field.id], dataset_version: self.dataset.manifest.version, - fragment_bitmap: if train { - match &self.fragments { - Some(fragment_ids) => Some(fragment_ids.iter().collect()), - None => Some(self.dataset.fragment_bitmap.as_ref().clone()), - } - } else { - // Empty bitmap for untrained indices - Some(roaring::RoaringBitmap::new()) - }, - index_details: Some(Arc::new(created_index.index_details)), + fragment_bitmap: Some(fragment_bitmap), + index_details: Some(Arc::new(index_details)), index_version: created_index.index_version as i32, created_at: Some(chrono::Utc::now()), base_id: None, diff --git a/rust/lance/src/index/vector.rs b/rust/lance/src/index/vector.rs index c4a6376a562..d2ff7b6035e 100644 --- a/rust/lance/src/index/vector.rs +++ b/rust/lance/src/index/vector.rs @@ -64,6 +64,7 @@ use tracing::instrument; use utils::get_vector_type; use uuid::Uuid; +use self::details::with_physical_fragment_bitmap; use super::{DatasetIndexExt, DatasetIndexInternalExt, IndexParams, pb}; use crate::dataset::index::dataset_format_version; use crate::dataset::transaction::{Operation, Transaction}; @@ -1941,15 +1942,22 @@ pub async fn initialize_vector_index( )) })?; - let fragment_bitmap = Some(target_dataset.fragment_bitmap.as_ref().clone()); + let fragment_bitmap = target_dataset.fragment_bitmap.as_ref().clone(); + let index_details = source_index + .index_details + .as_deref() + .cloned() + .map(|details| with_physical_fragment_bitmap(details, Some(&fragment_bitmap))) + .transpose()? + .map(Arc::new); let new_idx = IndexMetadata { uuid: new_uuid, name: source_index.name.clone(), fields: vec![field.id], dataset_version: target_dataset.manifest.version, - fragment_bitmap, - index_details: source_index.index_details.clone(), + fragment_bitmap: Some(fragment_bitmap), + index_details, index_version: source_index.index_version, created_at: Some(chrono::Utc::now()), base_id: None, diff --git a/rust/lance/src/index/vector/details.rs b/rust/lance/src/index/vector/details.rs index 3ccfd6e0e1e..8e9f6985de5 100644 --- a/rust/lance/src/index/vector/details.rs +++ b/rust/lance/src/index/vector/details.rs @@ -24,6 +24,7 @@ use lance_io::traits::Reader; use lance_io::utils::{CachedFileSize, read_last_block, read_version}; use lance_linalg::distance::DistanceType; use lance_table::format::IndexMetadata; +use roaring::RoaringBitmap; use serde::Serialize; use lance_index::vector::bq::{RQBuildParams, RQRotationType}; @@ -37,6 +38,67 @@ use crate::dataset::Dataset; use crate::index::open_index_proto; use crate::{Error, Result}; +/// Return the immutable physical fragment coverage recorded for a vector segment. +/// +/// Older segments do not carry this metadata and return `None` so callers can +/// choose a conservative compatibility path. +pub fn physical_fragment_bitmap(index: &IndexMetadata) -> Option { + let details = index.index_details.as_deref()?; + if !details.type_url.ends_with("VectorIndexDetails") { + return None; + } + let details = details.to_msg::().ok()?; + let encoded = details.physical_fragment_bitmap.as_deref()?; + let mut encoded = encoded; + RoaringBitmap::deserialize_from(&mut encoded).ok() +} + +/// Store immutable physical fragment coverage in vector index details. +/// +/// Passing `None` explicitly marks the coverage as unknown, which is required +/// when any input to a segment merge predates this metadata. +pub fn with_physical_fragment_bitmap( + details: prost_types::Any, + physical_fragments: Option<&RoaringBitmap>, +) -> Result { + if !details.type_url.ends_with("VectorIndexDetails") { + return Ok(details); + } + + let mut vector_details = details.to_msg::().map_err(|error| { + Error::index(format!( + "Failed to deserialize VectorIndexDetails while recording physical fragment coverage: {error}" + )) + })?; + vector_details.physical_fragment_bitmap = if let Some(bitmap) = physical_fragments { + let mut encoded = Vec::with_capacity(bitmap.serialized_size()); + bitmap.serialize_into(&mut encoded)?; + Some(encoded) + } else { + None + }; + prost_types::Any::from_msg(&vector_details).map_err(|error| { + Error::index(format!( + "Failed to serialize VectorIndexDetails with physical fragment coverage: {error}" + )) + }) +} + +/// Union physical coverage across merge inputs and newly owned fragments. +/// +/// A missing input provenance makes the result unknown instead of treating its +/// mutable ownership bitmap as physical coverage. +pub fn merged_physical_fragment_bitmap<'a>( + source_segments: impl IntoIterator, + newly_owned_fragments: &RoaringBitmap, +) -> Option { + let mut physical_fragments = newly_owned_fragments.clone(); + for source_segment in source_segments { + physical_fragments |= physical_fragment_bitmap(source_segment)?; + } + Some(physical_fragments) +} + // Private structs for JSON serialization of VectorIndexDetails. // Changes to field names or structure are backwards-incompatible for users // parsing the JSON output of describe_indices(). See snapshot tests below. @@ -162,6 +224,7 @@ pub fn vector_index_details(params: &VectorIndexParams) -> prost_types::Any { hnsw_index_config, compression, runtime_hints, + physical_fragment_bitmap: None, }; prost_types::Any::from_msg(&details).unwrap() } @@ -332,8 +395,8 @@ pub fn vector_params_from_details(details: &prost_types::Any) -> Option Option { let index_details = index.index_details.as_ref()?; - // Empty value bytes indicates legacy index that needs to be opened for details - if index_details.value.is_empty() { + // Physical coverage alone does not describe how to interpret the index. + if is_empty_vector_details(index_details) { return None; } @@ -345,10 +408,21 @@ pub fn metric_type_from_index_metadata(index: &IndexMetadata) -> Option bool { - details.value.is_empty() + if details.value.is_empty() { + return true; + } + details.to_msg::().is_ok_and(|details| { + details.metric_type == VectorMetricType::L2 as i32 + && details.target_partition_size == 0 + && details.hnsw_index_config.is_none() + && details.compression.is_none() + && details.runtime_hints.is_empty() + }) } /// Returns true if this is a vector index whose details need to be inferred from disk. @@ -363,7 +437,7 @@ pub fn needs_vector_details_inference( schema: &lance_core::datatypes::Schema, ) -> bool { match &index.index_details { - Some(d) => d.type_url.ends_with("VectorIndexDetails") && d.value.is_empty(), + Some(d) => d.type_url.ends_with("VectorIndexDetails") && is_empty_vector_details(d), None => index.fields.iter().any(|&field_id| { schema .field_by_id(field_id) @@ -406,7 +480,18 @@ pub async fn infer_missing_vector_details(dataset: &Dataset, indices: &mut [Inde .collect(); for index in indices.iter_mut() { if let Some(details) = inferred.get(&index.name) { - index.index_details = Some(details.clone()); + let physical_fragments = physical_fragment_bitmap(index); + match with_physical_fragment_bitmap( + details.as_ref().clone(), + physical_fragments.as_ref(), + ) { + Ok(details) => index.index_details = Some(Arc::new(details)), + Err(error) => tracing::warn!( + "Could not preserve vector index physical coverage for {}: {}", + index.name, + error + ), + } } } } @@ -558,6 +643,7 @@ fn convert_legacy_proto_to_details(proto: &pb::Index) -> Result no inference needed. diff --git a/rust/lance/src/index/vector/ivf.rs b/rust/lance/src/index/vector/ivf.rs index 2b4b602aa3d..0e5a385e569 100644 --- a/rust/lance/src/index/vector/ivf.rs +++ b/rust/lance/src/index/vector/ivf.rs @@ -5,6 +5,7 @@ use super::{ LogicalIvfView, derive_hnsw_params, + details::{merged_physical_fragment_bitmap, with_physical_fragment_bitmap}, pq::{PQIndex, build_pq_model}, utils::{filter_finite_training_data, maybe_sample_training_data}, }; @@ -2387,6 +2388,7 @@ pub(crate) async fn merge_segments_with_progress( })?; fragment_bitmap |= source_fragment_bitmap.clone(); } + let physical_fragment_bitmap = merged_physical_fragment_bitmap(&segments, &fragment_bitmap); let index_version = infer_source_index_version(&segments)?; let segment_uuid = Uuid::new_v4(); @@ -2404,7 +2406,10 @@ pub(crate) async fn merge_segments_with_progress( merged_segment = TableIndexMetadata { uuid: segment_uuid, fragment_bitmap: Some(fragment_bitmap), - index_details: Some(Arc::new(crate::index::vector_index_details_default())), + index_details: Some(Arc::new(with_physical_fragment_bitmap( + crate::index::vector_index_details_default(), + physical_fragment_bitmap.as_ref(), + )?)), index_version, created_at: Some(chrono::Utc::now()), base_id: None, diff --git a/rust/lance/src/index/vector/ivf/v2.rs b/rust/lance/src/index/vector/ivf/v2.rs index 63247cdaca4..0bd023eade7 100644 --- a/rust/lance/src/index/vector/ivf/v2.rs +++ b/rust/lance/src/index/vector/ivf/v2.rs @@ -3952,6 +3952,11 @@ mod tests { ); let expected_rows = fragments[0].physical_rows().await.unwrap() as u64 + fragments[1].physical_rows().await.unwrap() as u64; + let expected_physical_fragments = fragments + .iter() + .take(2) + .map(|fragment| fragment.id() as u32) + .collect(); let (ivf_params, pq_params) = prepare_global_ivf_pq(&dataset, "vector").await; let params = VectorIndexParams::with_ivf_pq_params(DistanceType::L2, ivf_params, pq_params); @@ -3977,6 +3982,11 @@ mod tests { ) .await .unwrap(); + assert_eq!( + crate::index::vector::details::physical_fragment_bitmap(&merged_segment), + Some(expected_physical_fragments), + "a merged segment must retain every input's physical fragment coverage" + ); dataset .commit_existing_index_segments(INDEX_NAME, "vector", vec![merged_segment]) .await diff --git a/rust/lance/src/io/exec/knn.rs b/rust/lance/src/io/exec/knn.rs index f952d5508f8..a40196f706a 100644 --- a/rust/lance/src/io/exec/knn.rs +++ b/rust/lance/src/io/exec/knn.rs @@ -63,9 +63,10 @@ use uuid::Uuid; use lance_select::RowAddrMask; use crate::dataset::Dataset; +use crate::index::DatasetIndexInternalExt; use crate::index::prefilter::{DatasetPreFilter, FilterLoader}; +use crate::index::vector::details::physical_fragment_bitmap; use crate::index::vector::utils::{get_vector_type, validate_distance_type_for}; -use crate::index::{DatasetIndexInternalExt, collect_subtree_field_ids, fragment_field_paths}; use crate::{Error, Result}; use lance_arrow::*; @@ -1596,9 +1597,8 @@ impl PreFilter for SegmentPreFilter { } /// Return whether a segment can still physically contain rows from a fragment it no -/// longer owns. Append-only deltas have disjoint bitmaps too, so the bitmap alone is -/// not enough: compare indexed data files with the snapshot used to build the segment. -async fn segment_needs_ownership_filter(dataset: &Dataset, index: &IndexMetadata) -> bool { +/// longer owns. +fn segment_needs_ownership_filter(dataset: &Dataset, index: &IndexMetadata) -> bool { let Some(owned_fragments) = index.fragment_bitmap.as_ref() else { return false; }; @@ -1610,45 +1610,14 @@ async fn segment_needs_ownership_filter(dataset: &Dataset, index: &IndexMetadata return false; } - let historical = match dataset.checkout_version(index.dataset_version).await { - Ok(historical) => historical, - Err(error) => { - // Old manifests may have been cleaned up. Restrict conservatively so a - // missing historical snapshot cannot make stale index rows visible again. - log::debug!( - "Could not inspect dataset version {} for index segment {}: {}. Applying the segment ownership filter conservatively.", - index.dataset_version, - index.uuid, - error - ); - return true; - } + let Some(physical_fragments) = physical_fragment_bitmap(index) else { + // Legacy vector segments do not record immutable physical coverage. A + // missing current fragment may be append-only or may still have stale + // rows in this segment, so filter conservatively for correctness. + return true; }; - - let mut indexed_field_ids = HashSet::new(); - for field_id in &index.fields { - if let Some(field) = dataset.schema().field_by_id(*field_id) { - collect_subtree_field_ids(field, &mut indexed_field_ids); - } - } - let current_fragments = dataset - .fragments() - .iter() - .map(|fragment| (fragment.id as u32, fragment)) - .collect::>(); - - historical.fragments().iter().any(|historical_fragment| { - let fragment_id = historical_fragment.id as u32; - if owned_fragments.contains(fragment_id) { - return false; - } - current_fragments - .get(&fragment_id) - .is_some_and(|current_fragment| { - fragment_field_paths(historical_fragment, &indexed_field_ids) - != fragment_field_paths(current_fragment, &indexed_field_ids) - }) - }) + let unowned_physical_fragments = physical_fragments - owned_fragments; + unowned_physical_fragments.intersection_len(&dataset.fragment_bitmap) > 0 } async fn prefilter_for_segment( @@ -1656,7 +1625,7 @@ async fn prefilter_for_segment( index: &IndexMetadata, base: Arc, ) -> Result> { - if !segment_needs_ownership_filter(dataset.as_ref(), index).await { + if !segment_needs_ownership_filter(dataset.as_ref(), index) { return Ok(base); } @@ -3000,26 +2969,27 @@ mod tests { ) .into_reader_rows(RowCount::from(20), BatchCount::from(1)); assert_eq!(second.schema(), first_schema); - dataset = Dataset::write( - second, - "memory://", - Some(WriteParams { - mode: WriteMode::Append, - ..Default::default() - }), - ) - .await - .unwrap(); + dataset.append(second, None).await.unwrap(); let appended_fragments = dataset.fragment_bitmap.as_ref() - &first_fragments; let dataset = Arc::new(dataset); + let old_details = crate::index::vector::details::with_physical_fragment_bitmap( + crate::index::vector_index_details_default(), + Some(&first_fragments), + ) + .unwrap(); + let new_details = crate::index::vector::details::with_physical_fragment_bitmap( + crate::index::vector_index_details_default(), + Some(&appended_fragments), + ) + .unwrap(); let old_segment = IndexMetadata { uuid: Uuid::new_v4(), fields: vec![field_id], name: "vector_idx".to_string(), dataset_version: first_version, - fragment_bitmap: Some(first_fragments), - index_details: None, + fragment_bitmap: Some(first_fragments.clone()), + index_details: Some(Arc::new(old_details)), index_version: 0, created_at: None, base_id: None, @@ -3030,17 +3000,44 @@ mod tests { fields: vec![field_id], name: "vector_idx".to_string(), dataset_version: dataset.manifest.version, - fragment_bitmap: Some(appended_fragments), - index_details: None, + fragment_bitmap: Some(appended_fragments.clone()), + index_details: Some(Arc::new(new_details)), index_version: 0, created_at: None, base_id: None, files: None, }; assert!( - !segment_needs_ownership_filter(dataset.as_ref(), &old_segment).await, + !segment_needs_ownership_filter(dataset.as_ref(), &old_segment), "append-only data must not look like an in-place fragment update" ); + + let merged_physical_fragments = &first_fragments | &appended_fragments; + let merged_details = crate::index::vector::details::with_physical_fragment_bitmap( + crate::index::vector_index_details_default(), + Some(&merged_physical_fragments), + ) + .unwrap(); + let merged_with_pruned_later_fragment = IndexMetadata { + index_details: Some(Arc::new(merged_details)), + ..old_segment.clone() + }; + assert_eq!( + physical_fragment_bitmap(&merged_with_pruned_later_fragment), + Some(merged_physical_fragments.clone()) + ); + assert!( + !appended_fragments.is_empty(), + "the append fixture must create at least one new fragment" + ); + assert!( + appended_fragments.intersection_len(&dataset.fragment_bitmap) > 0, + "the appended fragments must remain live in the current dataset" + ); + assert!( + segment_needs_ownership_filter(dataset.as_ref(), &merged_with_pruned_later_fragment), + "a merged segment must retain the physical coverage of every input" + ); let base = Arc::new(DatasetPreFilter::new( dataset.clone(), &[old_segment.clone(), new_segment], From 6de9820644d13b2409fab1676f7cde036e562c4b Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Thu, 6 Aug 2026 16:47:57 +0000 Subject: [PATCH 3/7] test(index): account for merged physical coverage --- rust/lance/src/index/create.rs | 29 ++++++++++++++++++++++++----- 1 file changed, 24 insertions(+), 5 deletions(-) diff --git a/rust/lance/src/index/create.rs b/rust/lance/src/index/create.rs index e954cb88d11..c294be10700 100644 --- a/rust/lance/src/index/create.rs +++ b/rust/lance/src/index/create.rs @@ -1018,7 +1018,7 @@ impl<'a> IntoFuture for CreateIndexBuilder<'a> { mod tests { use super::*; use crate::dataset::{WriteMode, WriteParams}; - use crate::index::{DatasetIndexExt, IndexSegment}; + use crate::index::{DatasetIndexExt, IndexSegment, vector::details::physical_fragment_bitmap}; use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; use arrow::datatypes::{Float32Type, Int32Type, Int64Type}; use arrow_array::cast::AsArray; @@ -1029,7 +1029,6 @@ mod tests { use lance_arrow::FixedSizeListArrayExt; use lance_core::utils::{address::RowAddress, tempfile::TempStrDir}; use lance_datagen::{self, gen_batch}; - use lance_index::optimize::OptimizeOptions; use lance_index::progress::IndexBuildProgress; use lance_index::scalar::{ BloomFilterQuery, FullTextSearchQuery, SargableQuery, SearchResult, @@ -1038,6 +1037,7 @@ mod tests { use lance_index::vector::hnsw::builder::HnswBuildParams; use lance_index::vector::ivf::IvfBuildParams; use lance_index::vector::kmeans::{KMeansParams, train_kmeans}; + use lance_index::{optimize::OptimizeOptions, pb::VectorIndexDetails}; use lance_linalg::distance::{DistanceType, MetricType}; use roaring::RoaringBitmap; use std::{collections::BTreeSet, ops::Bound, sync::Arc}; @@ -1351,6 +1351,17 @@ mod tests { .unwrap() } + fn vector_details_without_physical_coverage(index: &IndexMetadata) -> VectorIndexDetails { + let mut details = index + .index_details + .as_deref() + .expect("vector index should have details") + .to_msg::() + .expect("vector index details should decode"); + details.physical_fragment_bitmap = None; + details + } + #[tokio::test] async fn test_get_frags_from_ordered_ids_accepts_unsorted_duplicates() { let tmpdir = TempStrDir::default(); @@ -2541,9 +2552,11 @@ mod tests { .iter() .flat_map(|segment| segment.fragment_bitmap.as_ref().unwrap().iter()) .collect::(); - let expected_merged_details = compatible_tail.last().unwrap().index_details.clone(); + let expected_merged_details = + vector_details_without_physical_coverage(compatible_tail.last().unwrap()); assert_ne!( - independent_segment.index_details, expected_merged_details, + vector_details_without_physical_coverage(&independent_segment), + expected_merged_details, "test setup must distinguish base and suffix metadata" ); let mut segments = vec![independent_segment]; @@ -2581,9 +2594,15 @@ mod tests { &expected_merged_fragments ); assert_eq!( - merged.index_details, expected_merged_details, + vector_details_without_physical_coverage(merged), + expected_merged_details, "merged metadata must come from the selected suffix, not an incompatible base segment" ); + assert_eq!( + physical_fragment_bitmap(merged), + Some(expected_merged_fragments), + "merged physical coverage must include every selected suffix segment" + ); } #[tokio::test] From 2e0f73b47dbde87e69982e6abd9d3496543258f5 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Sat, 8 Aug 2026 12:27:48 +0000 Subject: [PATCH 4/7] fix(index): filter ownership during vector segment merges --- python/python/tests/test_vector_index.py | 26 ++- .../src/vector/distributed/index_merger.rs | 205 +++++++++++++++--- rust/lance/src/index.rs | 101 +++++++-- rust/lance/src/index/append.rs | 74 ++++++- rust/lance/src/index/vector/details.rs | 28 ++- rust/lance/src/index/vector/ivf.rs | 96 ++++++-- 6 files changed, 436 insertions(+), 94 deletions(-) diff --git a/python/python/tests/test_vector_index.py b/python/python/tests/test_vector_index.py index 5b95e898363..6c0b29a7de1 100644 --- a/python/python/tests/test_vector_index.py +++ b/python/python/tests/test_vector_index.py @@ -2012,17 +2012,23 @@ def table(ids, value): dataset.optimize.optimize_indices(num_indices_to_merge=0) dataset = lance.dataset(dataset.uri) - result = dataset.to_table( - columns=["id"], - nearest={ - "column": "vector", - "q": np.zeros(ndim, dtype=np.float32), - "k": 5, - }, - ) + def assert_current_nearest_rows(): + result = dataset.to_table( + columns=["id"], + nearest={ + "column": "vector", + "q": np.zeros(ndim, dtype=np.float32), + "k": 5, + }, + ) + + assert all(row_id < 20 for row_id in result["id"].to_pylist()) + assert result["_distance"].to_pylist() == pytest.approx([4.0] * 5) - assert all(row_id < 20 for row_id in result["id"].to_pylist()) - assert result["_distance"].to_pylist() == pytest.approx([4.0] * 5) + assert_current_nearest_rows() + dataset.optimize.optimize_indices(num_indices_to_merge=2) + dataset = lance.dataset(dataset.uri) + assert_current_nearest_rows() @pytest.mark.skip(reason="retrain is deprecated") diff --git a/rust/lance-index/src/vector/distributed/index_merger.rs b/rust/lance-index/src/vector/distributed/index_merger.rs index 2b0345b4829..eb91ac46432 100755 --- a/rust/lance-index/src/vector/distributed/index_merger.rs +++ b/rust/lance-index/src/vector/distributed/index_merger.rs @@ -9,16 +9,17 @@ use crate::vector::shared::partition_merger::{ }; use arrow::{compute::concat_batches, datatypes::Float32Type}; use arrow_array::cast::AsArray; -use arrow_array::types::UInt8Type; +use arrow_array::types::{UInt8Type, UInt64Type}; use arrow_array::{Array, FixedSizeListArray, RecordBatch}; use futures::StreamExt as _; use lance_arrow::{FixedSizeListArrayExt, RecordBatchExt}; -use lance_core::{Error, ROW_ID_FIELD, Result}; +use lance_core::{Error, ROW_ID, ROW_ID_FIELD, Result}; use std::ops::Range; use std::sync::Arc; use crate::IndexMetadata as IndexMetaSchema; use crate::pb; +use crate::scalar::OldIndexDataFilter; use crate::vector::bq::storage::{ RABIT_CODE_COLUMN, RABIT_METADATA_KEY, RabitQuantizationMetadata, RabitQueryEstimator, pack_codes, rabit_binary_code_field, rabit_ex_code_field, @@ -525,6 +526,7 @@ struct ShardInfo { lengths: Vec, partition_offsets: Vec, total_rows: usize, + row_filter: Option>, } #[derive(Debug)] @@ -534,6 +536,7 @@ struct ShardWindowReadJob { window_total_rows: usize, start_offset: usize, end_offset: usize, + row_filter: Option>, } #[derive(Debug)] @@ -652,6 +655,7 @@ async fn read_partition_window( window_total_rows, start_offset, end_offset, + row_filter: shard.row_filter.clone(), } }) .collect(); @@ -739,7 +743,13 @@ async fn read_shard_window_partitions( } let to_take = std::cmp::min(remaining, rb.num_rows() - consumed); - per_partition_batches[rel_partition].push(rb.slice(consumed, to_take)); + let mut partition_batch = rb.slice(consumed, to_take); + if let Some(row_filter) = shard_job.row_filter.as_deref() { + partition_batch = filter_batch_to_owned_rows(&partition_batch, row_filter)?; + } + if partition_batch.num_rows() > 0 { + per_partition_batches[rel_partition].push(partition_batch); + } consumed += to_take; remaining -= to_take; } @@ -762,6 +772,19 @@ async fn read_shard_window_partitions( Ok(per_partition_batches) } +fn filter_batch_to_owned_rows( + batch: &RecordBatch, + row_filter: &OldIndexDataFilter, +) -> Result { + let row_ids = batch + .column_by_name(ROW_ID) + .ok_or_else(|| Error::index(format!("Column {ROW_ID} missing in auxiliary shard")))? + .as_primitive_opt::() + .ok_or_else(|| Error::index(format!("Column {ROW_ID} is not UInt64 in auxiliary shard")))?; + let keep = row_filter.filter_row_ids(row_ids); + Ok(arrow::compute::filter_record_batch(batch, &keep)?) +} + /// Merge the selected segment auxiliary files into `target_dir`. /// /// This is the storage merge kernel for vector segment build. Callers choose @@ -778,6 +801,42 @@ pub async fn merge_partial_vector_auxiliary_files( aux_paths: &[object_store::path::Path], target_dir: &object_store::path::Path, progress: Arc, +) -> Result { + merge_partial_vector_auxiliary_files_inner(object_store, aux_paths, target_dir, None, progress) + .await +} + +/// Merge auxiliary files while retaining only rows owned by each source segment. +pub async fn merge_partial_vector_auxiliary_files_with_row_filters( + object_store: &lance_io::object_store::ObjectStore, + aux_paths: &[object_store::path::Path], + target_dir: &object_store::path::Path, + row_filters: &[OldIndexDataFilter], + progress: Arc, +) -> Result { + if aux_paths.len() != row_filters.len() { + return Err(Error::invalid_input(format!( + "Expected one row filter per auxiliary file, got {} files and {} filters", + aux_paths.len(), + row_filters.len() + ))); + } + merge_partial_vector_auxiliary_files_inner( + object_store, + aux_paths, + target_dir, + Some(row_filters), + progress, + ) + .await +} + +async fn merge_partial_vector_auxiliary_files_inner( + object_store: &lance_io::object_store::ObjectStore, + aux_paths: &[object_store::path::Path], + target_dir: &object_store::path::Path, + row_filters: Option<&[OldIndexDataFilter]>, + progress: Arc, ) -> Result { if aux_paths.is_empty() { return Err(Error::index( @@ -1454,6 +1513,7 @@ pub async fn merge_partial_vector_auxiliary_files( lengths, partition_offsets, total_rows: running_offset, + row_filter: row_filters.map(|filters| Arc::new(filters[idx].clone())), }); progress .stage_progress("read_shard_metadata", idx as u64 + 1) @@ -1480,6 +1540,7 @@ pub async fn merge_partial_vector_auxiliary_files( .stage_start("merge_partitions", Some(total_rows), "rows") .await?; let mut merged_rows = 0u64; + let mut merged_lengths = vec![0u32; nlist]; match idx_type_final { SupportedIvfIndexType::IvfPq | SupportedIvfIndexType::IvfHnswPq => { @@ -1495,22 +1556,20 @@ pub async fn merge_partial_vector_auxiliary_files( ); while let Some((pid, batches)) = shard_merge_reader.next_partition().await? { - if accumulated_lengths[pid] == 0 { + let partition_len = batches.iter().map(RecordBatch::num_rows).sum::(); + if partition_len == 0 { continue; } - if batches.is_empty() { - return Err(Error::index(format!( - "No merged batches found for non-empty partition {}", - pid - ))); - } let schema = batches[0].schema(); let partition_batch = concat_batches(&schema, batches.iter())?; if let Some(w) = v2w_opt.as_mut() { write_partition_rows_pq_transposed(w, partition_batch).await?; } - merged_rows = merged_rows.saturating_add(accumulated_lengths[pid] as u64); + merged_lengths[pid] = u32::try_from(partition_len).map_err(|_| { + Error::index(format!("Merged partition {pid} exceeds u32 row capacity")) + })?; + merged_rows = merged_rows.saturating_add(partition_len as u64); progress .stage_progress("merge_partitions", merged_rows) .await?; @@ -1527,15 +1586,10 @@ pub async fn merge_partial_vector_auxiliary_files( ); while let Some((pid, batches)) = shard_merge_reader.next_partition().await? { - if accumulated_lengths[pid] == 0 { + let partition_len = batches.iter().map(RecordBatch::num_rows).sum::(); + if partition_len == 0 { continue; } - if batches.is_empty() { - return Err(Error::index(format!( - "No merged batches found for non-empty partition {}", - pid - ))); - } // Shards written by older lance versions carry sequential ex // codes; normalize every batch to the blocked layout before @@ -1561,30 +1615,38 @@ pub async fn merge_partial_vector_auxiliary_files( if let Some(w) = v2w_opt.as_mut() { write_partition_rows_rq_packed(w, partition_batch).await?; } - merged_rows = merged_rows.saturating_add(accumulated_lengths[pid] as u64); + merged_lengths[pid] = u32::try_from(partition_len).map_err(|_| { + Error::index(format!("Merged partition {pid} exceeds u32 row capacity")) + })?; + merged_rows = merged_rows.saturating_add(partition_len as u64); progress .stage_progress("merge_partitions", merged_rows) .await?; } } _ => { - for (pid, total_part_len) in accumulated_lengths.iter().copied().enumerate().take(nlist) - { - for shard in shard_infos.iter() { - let part_len = shard.lengths[pid] as usize; - if part_len == 0 { - continue; - } - let offset = shard.partition_offsets[pid]; - if let Some(w) = v2w_opt.as_mut() { - write_partition_rows(shard.reader.as_ref(), w, offset..offset + part_len) - .await?; - } - } - if total_part_len == 0 { + let partition_window_size = *PARTITION_WINDOW_SIZE; + let prefetch_window_count = *PARTITION_PREFETCH_WINDOW_COUNT; + let mut shard_merge_reader = ShardMergeReader::new( + shard_infos, + nlist, + partition_window_size, + prefetch_window_count, + ); + while let Some((pid, batches)) = shard_merge_reader.next_partition().await? { + let partition_len = batches.iter().map(RecordBatch::num_rows).sum::(); + if partition_len == 0 { continue; } - merged_rows = merged_rows.saturating_add(total_part_len as u64); + if let Some(w) = v2w_opt.as_mut() { + for batch in batches { + w.write_batch(&batch).await?; + } + } + merged_lengths[pid] = u32::try_from(partition_len).map_err(|_| { + Error::index(format!("Merged partition {pid} exceeds u32 row capacity")) + })?; + merged_rows = merged_rows.saturating_add(partition_len as u64); progress .stage_progress("merge_partitions", merged_rows) .await?; @@ -1603,7 +1665,7 @@ pub async fn merge_partial_vector_auxiliary_files( } else { IvfStorageModel::empty() }; - for len in accumulated_lengths.iter() { + for len in merged_lengths.iter() { ivf_model.add_partition(*len); } let dt2 = distance_type.ok_or_else(|| Error::index("Distance type missing".to_string()))?; @@ -1980,6 +2042,79 @@ mod tests { assert_eq!(total_rows, expected_total); } + #[tokio::test] + async fn test_merge_ivf_flat_filters_each_source_by_ownership() { + let object_store = ObjectStore::memory(); + let index_dir = Path::from("index/uuid"); + let aux0 = index_dir + .clone() + .join("stale") + .join(INDEX_AUXILIARY_FILE_NAME); + let aux1 = index_dir + .clone() + .join("fresh") + .join(INDEX_AUXILIARY_FILE_NAME); + let lengths = vec![2_u32, 1_u32]; + + write_flat_partial_aux(&object_store, &aux0, 2, &lengths, 0, DistanceType::L2) + .await + .unwrap(); + write_flat_partial_aux(&object_store, &aux1, 2, &lengths, 100, DistanceType::L2) + .await + .unwrap(); + + merge_partial_vector_auxiliary_files_with_row_filters( + &object_store, + &[aux0, aux1], + &index_dir, + &[ + OldIndexDataFilter::Fragments { + to_keep: roaring::RoaringBitmap::new(), + to_remove: roaring::RoaringBitmap::new(), + }, + OldIndexDataFilter::RowIds(lance_select::RowAddrTreeMap::from_iter(100_u64..103)), + ], + Arc::new(RecordingProgress::default()), + ) + .await + .unwrap(); + + let aux_out = index_dir.join(INDEX_AUXILIARY_FILE_NAME); + let sched = ScanScheduler::new( + Arc::new(object_store.clone()), + SchedulerConfig::max_bandwidth(&object_store), + ); + let reader = V2Reader::try_open( + sched + .open_file(&aux_out, &CachedFileSize::unknown()) + .await + .unwrap(), + None, + Arc::default(), + &lance_core::cache::LanceCache::no_cache(), + V2ReaderOptions::default(), + ) + .await + .unwrap(); + let merged_ivf = try_read_ivf_proto(&reader).await.unwrap().unwrap(); + assert_eq!(merged_ivf.lengths, lengths); + + let mut total_rows = 0; + let mut stream = reader + .read_stream( + lance_io::ReadBatchParams::RangeFull, + u32::MAX, + 4, + lance_encoding::decoder::FilterExpression::no_filter(), + ) + .await + .unwrap(); + while let Some(batch) = stream.next().await { + total_rows += batch.unwrap().num_rows(); + } + assert_eq!(total_rows, 3, "stale source rows must not be copied"); + } + #[tokio::test] async fn test_merge_distance_type_mismatch() { let object_store = ObjectStore::memory(); diff --git a/rust/lance/src/index.rs b/rust/lance/src/index.rs index 163ab3f1c40..9733b996cf1 100644 --- a/rust/lance/src/index.rs +++ b/rust/lance/src/index.rs @@ -132,6 +132,24 @@ fn validate_segment_metadata(index_name: &str, segments: &[IndexMetadata]) -> Re Ok(()) } +fn project_index_through_fragment_reuse( + index: &mut IndexMetadata, + frag_reuse_index: &FragReuseIndex, +) -> Result<()> { + let Some(fragment_bitmap) = index.fragment_bitmap.as_mut() else { + return Ok(()); + }; + frag_reuse_index.remap_fragment_bitmap(fragment_bitmap)?; + index.index_details = index + .index_details + .as_deref() + .cloned() + .map(|details| with_physical_fragment_bitmap(details, None)) + .transpose()? + .map(Arc::new); + Ok(()) +} + fn collect_subtree_field_ids(field: &Field, field_ids: &mut HashSet) { field_ids.insert(field.id); for child in &field.children { @@ -1553,9 +1571,7 @@ impl DatasetIndexExt for Dataset { .await?; let mut indices = indices.as_ref().clone(); for idx in indices.iter_mut() { - if let Some(bitmap) = idx.fragment_bitmap.as_mut() { - frag_reuse_index.remap_fragment_bitmap(bitmap)?; - } + project_index_through_fragment_reuse(idx, &frag_reuse_index)?; } Ok(Arc::new(indices)) } else { @@ -1639,12 +1655,7 @@ impl DatasetIndexExt for Dataset { }; let mut merged_segment = if all_vector { - crate::index::vector::ivf::merge_segments( - self.object_store.as_ref(), - &self.indices_dir(), - source_segments, - ) - .await? + crate::index::vector::ivf::merge_segments(self, source_segments).await? } else if all_inverted { crate::index::scalar::inverted::merge_segments(self, source_segments).await? } else if all_fmindex { @@ -1947,10 +1958,12 @@ impl DatasetIndexExt for Dataset { }; let last_idx = deltas.last().expect("Delta indices should not be empty"); - let physical_fragment_bitmap = merged_physical_fragment_bitmap( - res.removed_indices.iter().copied(), - &res.new_fragment_bitmap, - ); + let physical_fragment_bitmap = res.new_physical_fragment_bitmap.or_else(|| { + merged_physical_fragment_bitmap( + res.removed_indices.iter().copied(), + &res.new_fragment_bitmap, + ) + }); let new_index_details = with_physical_fragment_bitmap( res.new_index_details, physical_fragment_bitmap.as_ref(), @@ -3192,6 +3205,11 @@ mod tests { use lance_core::utils::tempfile::TempStrDir; use lance_datagen::gen_batch; use lance_datagen::{BatchCount, ByteCount, Dimension, RowCount, array}; + use lance_index::frag_reuse::{ + FragDigest, FragReuseGroup, FragReuseIndexDetails, FragReuseVersion, + }; + use lance_index::pb::VectorIndexDetails; + use lance_index::pb::vector_index_details::{Compression, FlatCompression}; use lance_index::pbold::{BTreeIndexDetails, InvertedIndexDetails}; use lance_index::scalar::bitmap::BITMAP_LOOKUP_NAME; use lance_index::scalar::inverted::query::{FtsQuery, PhraseQuery}; @@ -3214,6 +3232,63 @@ mod tests { use rstest::rstest; use std::collections::{HashMap, HashSet}; + #[test] + fn test_fragment_reuse_projection_invalidates_vector_physical_coverage() { + let physical_fragments = RoaringBitmap::from_iter([0_u32]); + let details = prost_types::Any::from_msg(&VectorIndexDetails { + compression: Some(Compression::Flat(FlatCompression {})), + ..Default::default() + }) + .unwrap(); + let details = with_physical_fragment_bitmap(details, Some(&physical_fragments)).unwrap(); + let mut index = IndexMetadata { + uuid: Uuid::new_v4(), + name: "vector_idx".to_string(), + fields: vec![0], + dataset_version: 1, + fragment_bitmap: Some(physical_fragments), + index_details: Some(Arc::new(details)), + index_version: 1, + created_at: None, + base_id: None, + files: None, + }; + let digest = |id| FragDigest { + id, + physical_rows: 1, + num_deleted_rows: 0, + }; + let frag_reuse_index = FragReuseIndex::new( + Uuid::new_v4(), + Vec::new(), + FragReuseIndexDetails { + versions: vec![FragReuseVersion { + dataset_version: 2, + groups: vec![FragReuseGroup { + changed_row_addrs: Vec::new(), + old_frags: vec![digest(0)], + new_frags: vec![digest(1)], + }], + }], + }, + ); + + project_index_through_fragment_reuse(&mut index, &frag_reuse_index).unwrap(); + + assert_eq!(index.fragment_bitmap, Some(RoaringBitmap::from_iter([1]))); + assert_eq!(vector::details::physical_fragment_bitmap(&index), None); + let remapped_details = index + .index_details + .as_deref() + .unwrap() + .to_msg::() + .unwrap(); + assert!(matches!( + remapped_details.compression, + Some(Compression::Flat(_)) + )); + } + async fn write_vector_segment_metadata( dataset: &Dataset, index_name: &str, diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 398d8ee8170..db27cd1732a 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -40,6 +40,7 @@ use crate::dataset::rowids::load_row_id_sequences; use crate::index::scalar::{ IndexDetails, fetch_index_details, load_fts_training_data, load_training_data, }; +use crate::index::vector::details::physical_fragment_bitmap; use crate::index::vector_index_details_default; #[derive(Debug, Clone)] @@ -47,6 +48,8 @@ pub struct IndexMergeResults<'a> { pub new_uuid: Uuid, pub removed_indices: Vec<&'a IndexMetadata>, pub new_fragment_bitmap: RoaringBitmap, + /// Exact physical coverage when the output was rebuilt from current rows. + pub new_physical_fragment_bitmap: Option, pub new_dataset_version: u64, pub new_index_version: i32, pub new_index_details: prost_types::Any, @@ -687,6 +690,21 @@ async fn build_fresh_vector_segment( .await } +fn segment_merge_requires_rebuild( + index: &IndexMetadata, + current_fragments: &RoaringBitmap, +) -> bool { + let Some(physical_fragments) = physical_fragment_bitmap(index) else { + return true; + }; + let owned_fragments = index + .effective_fragment_bitmap(current_fragments) + .or_else(|| index.fragment_bitmap.clone()) + .unwrap_or_default(); + let unowned_physical_fragments = &physical_fragments - &owned_fragments; + unowned_physical_fragments.intersection_len(current_fragments) > 0 +} + async fn scan_vector_fragments( dataset: &Dataset, field_path: &str, @@ -731,6 +749,7 @@ fn fresh_vector_segment_result<'a>( Ok(IndexMergeResults { new_uuid: segment.uuid, removed_indices, + new_physical_fragment_bitmap: Some(fragment_bitmap.clone()), new_fragment_bitmap: fragment_bitmap, new_dataset_version: segment.dataset_version, new_index_version: segment.index_version, @@ -943,6 +962,7 @@ pub async fn merge_indices_with_unindexed_frags<'a>( new_uuid, removed_indices: Vec::new(), new_fragment_bitmap: base_unindexed_bitmap, + new_physical_fragment_bitmap: None, new_dataset_version: dataset.manifest.version, new_index_version: index_type_for_segmented_optimize(reference_index.as_ref())? .version(), @@ -1001,6 +1021,26 @@ pub async fn merge_indices_with_unindexed_frags<'a>( field_path.clone(), vec![(selected_metadata, selected_index)], )?; + let new_fragment_bitmap = removed_segment + .effective_fragment_bitmap(&dataset.fragment_bitmap) + .or_else(|| removed_segment.fragment_bitmap.clone()) + .unwrap_or_default(); + if segment_merge_requires_rebuild(removed_segment, &dataset.fragment_bitmap) { + let segment = build_fresh_vector_segment( + dataset.as_ref(), + &selected_logical_index, + &field_path, + &new_fragment_bitmap, + options.progress.clone(), + ) + .await?; + return fresh_vector_segment_result( + segment, + &new_fragment_bitmap, + vec![removed_segment], + ) + .map(Some); + } let selected_ivf_view = selected_logical_index.as_ivf()?; let (new_uuid, indices_merged, files) = Box::pin(optimize_vector_indices( dataset.as_ref().clone(), @@ -1018,11 +1058,6 @@ pub async fn merge_indices_with_unindexed_frags<'a>( return Ok(None); } - let new_fragment_bitmap = removed_segment - .effective_fragment_bitmap(&dataset.fragment_bitmap) - .or_else(|| removed_segment.fragment_bitmap.clone()) - .unwrap_or_default(); - Ok(( new_uuid, vec![removed_segment], @@ -1056,6 +1091,33 @@ pub async fn merge_indices_with_unindexed_frags<'a>( .map(|(metadata, index)| (metadata.clone(), index.clone())) .collect(), )?; + let selected_segments = &old_indices[merge_start..]; + if selected_segments.iter().any(|segment| { + segment_merge_requires_rebuild(segment, &dataset.fragment_bitmap) + }) { + let mut fragment_bitmap = base_unindexed_bitmap.clone(); + for segment in selected_segments { + if let Some(effective) = + segment.effective_fragment_bitmap(&dataset.fragment_bitmap) + { + fragment_bitmap |= effective; + } + } + let segment = build_fresh_vector_segment( + dataset.as_ref(), + &merge_logical_index, + &field_path, + &fragment_bitmap, + options.progress.clone(), + ) + .await?; + return fresh_vector_segment_result( + segment, + &fragment_bitmap, + selected_segments.to_vec(), + ) + .map(Some); + } let merge_ivf_view = merge_logical_index.as_ivf()?; let new_data_stream = if unindexed.is_empty() { @@ -1218,6 +1280,7 @@ pub async fn merge_indices_with_unindexed_frags<'a>( new_uuid, removed_indices: old_indices.to_vec(), new_fragment_bitmap: dataset.fragment_bitmap.as_ref().clone(), + new_physical_fragment_bitmap: None, new_dataset_version: dataset.manifest.version, new_index_version: created_index.index_version as i32, new_index_details: created_index.index_details, @@ -1352,6 +1415,7 @@ pub async fn merge_indices_with_unindexed_frags<'a>( new_uuid, removed_indices, new_fragment_bitmap, + new_physical_fragment_bitmap: None, new_dataset_version, new_index_version: created_index.index_version as i32, new_index_details: created_index.index_details, diff --git a/rust/lance/src/index/vector/details.rs b/rust/lance/src/index/vector/details.rs index 8e9f6985de5..42738bec1a9 100644 --- a/rust/lance/src/index/vector/details.rs +++ b/rust/lance/src/index/vector/details.rs @@ -70,7 +70,9 @@ pub fn with_physical_fragment_bitmap( "Failed to deserialize VectorIndexDetails while recording physical fragment coverage: {error}" )) })?; - vector_details.physical_fragment_bitmap = if let Some(bitmap) = physical_fragments { + vector_details.physical_fragment_bitmap = if let Some(bitmap) = physical_fragments + && has_vector_index_configuration(&vector_details) + { let mut encoded = Vec::with_capacity(bitmap.serialized_size()); bitmap.serialize_into(&mut encoded)?; Some(encoded) @@ -412,17 +414,21 @@ pub fn metric_type_from_index_metadata(index: &IndexMetadata) -> Option bool { + details.metric_type != VectorMetricType::L2 as i32 + || details.target_partition_size != 0 + || details.hnsw_index_config.is_some() + || details.compression.is_some() + || !details.runtime_hints.is_empty() +} + fn is_empty_vector_details(details: &prost_types::Any) -> bool { if details.value.is_empty() { return true; } - details.to_msg::().is_ok_and(|details| { - details.metric_type == VectorMetricType::L2 as i32 - && details.target_partition_size == 0 - && details.hnsw_index_config.is_none() - && details.compression.is_none() - && details.runtime_hints.is_empty() - }) + details + .to_msg::() + .is_ok_and(|details| !has_vector_index_configuration(&details)) } /// Returns true if this is a vector index whose details need to be inferred from disk. @@ -1137,7 +1143,7 @@ mod tests { } #[test] - fn test_physical_coverage_does_not_suppress_details_inference() { + fn test_physical_coverage_preserves_empty_details_sentinel() { let schema = schema_with_vector_and_scalar(); let vec_id = schema.field("vec").unwrap().id; let physical_fragments = RoaringBitmap::from_iter([1, 3]); @@ -1146,9 +1152,9 @@ mod tests { Some(&physical_fragments), ) .unwrap(); + assert!(details.value.is_empty()); let index = index_over_field(vec_id, Some(details)); - - assert_eq!(physical_fragment_bitmap(&index), Some(physical_fragments)); + assert_eq!(physical_fragment_bitmap(&index), None); assert!(needs_vector_details_inference(&index, &schema)); assert_eq!(metric_type_from_index_metadata(&index), None); assert_eq!( diff --git a/rust/lance/src/index/vector/ivf.rs b/rust/lance/src/index/vector/ivf.rs index 0e5a385e569..f3607275bea 100644 --- a/rust/lance/src/index/vector/ivf.rs +++ b/rust/lance/src/index/vector/ivf.rs @@ -5,7 +5,7 @@ use super::{ LogicalIvfView, derive_hnsw_params, - details::{merged_physical_fragment_bitmap, with_physical_fragment_bitmap}, + details::with_physical_fragment_bitmap, pq::{PQIndex, build_pq_model}, utils::{filter_finite_training_data, maybe_sample_training_data}, }; @@ -2349,14 +2349,35 @@ async fn write_ivf_hnsw_file( /// Merge one caller-defined group of source segments into a single segment. pub(crate) async fn merge_segments( - object_store: &ObjectStore, - indices_dir: &Path, + dataset: &Dataset, segments: Vec, ) -> Result { - merge_segments_with_progress( - object_store, - indices_dir, + let mut row_filters = Vec::with_capacity(segments.len()); + let no_deleted_fragments = RoaringBitmap::new(); + for segment in &segments { + let owned_fragments = segment.fragment_bitmap.as_ref().ok_or_else(|| { + Error::index(format!( + "Segment '{}' is missing fragment coverage", + segment.uuid + )) + })?; + row_filters.push( + crate::index::append::build_old_data_filter( + dataset, + owned_fragments, + &no_deleted_fragments, + ) + .await? + .ok_or_else(|| { + Error::internal("Vector segment ownership filter is missing".to_string()) + })?, + ); + } + merge_segments_with_row_filters( + dataset.object_store.as_ref(), + &dataset.indices_dir(), segments, + row_filters, lance_index::progress::noop_progress(), ) .await @@ -2364,11 +2385,38 @@ pub(crate) async fn merge_segments( /// Merge one caller-defined group of source segments into a single segment and /// report progress through the provided callback. +#[cfg(test)] pub(crate) async fn merge_segments_with_progress( object_store: &ObjectStore, indices_dir: &Path, segments: Vec, progress: Arc, +) -> Result { + let row_filters = segments + .iter() + .map(|segment| { + let to_keep = segment.fragment_bitmap.clone().ok_or_else(|| { + Error::index(format!( + "Segment '{}' is missing fragment coverage", + segment.uuid + )) + })?; + Ok(lance_index::scalar::OldIndexDataFilter::Fragments { + to_keep, + to_remove: RoaringBitmap::new(), + }) + }) + .collect::>>()?; + merge_segments_with_row_filters(object_store, indices_dir, segments, row_filters, progress) + .await +} + +async fn merge_segments_with_row_filters( + object_store: &ObjectStore, + indices_dir: &Path, + segments: Vec, + row_filters: Vec, + progress: Arc, ) -> Result { if segments.is_empty() { return Err(Error::index("No segment metadata was provided".to_string())); @@ -2388,7 +2436,16 @@ pub(crate) async fn merge_segments_with_progress( })?; fragment_bitmap |= source_fragment_bitmap.clone(); } - let physical_fragment_bitmap = merged_physical_fragment_bitmap(&segments, &fragment_bitmap); + let mut index_details = crate::index::vector_index_details_default(); + for segment in &segments { + if let Some(details) = segment.index_details.as_deref() { + let details = with_physical_fragment_bitmap(details.clone(), None)?; + if !details.value.is_empty() { + index_details = details; + break; + } + } + } let index_version = infer_source_index_version(&segments)?; let segment_uuid = Uuid::new_v4(); @@ -2398,18 +2455,17 @@ pub(crate) async fn merge_segments_with_progress( indices_dir, &final_dir, &segments, + &row_filters, None, progress, ) .await?; + let index_details = with_physical_fragment_bitmap(index_details, Some(&fragment_bitmap))?; merged_segment = TableIndexMetadata { uuid: segment_uuid, fragment_bitmap: Some(fragment_bitmap), - index_details: Some(Arc::new(with_physical_fragment_bitmap( - crate::index::vector_index_details_default(), - physical_fragment_bitmap.as_ref(), - )?)), + index_details: Some(Arc::new(index_details)), index_version, created_at: Some(chrono::Utc::now()), base_id: None, @@ -2429,6 +2485,7 @@ async fn merge_segments_to_dir( indices_dir: &Path, final_dir: &Path, segments: &[TableIndexMetadata], + row_filters: &[lance_index::scalar::OldIndexDataFilter], _requested_index_type: Option, progress: Arc, ) -> Result> { @@ -2457,15 +2514,14 @@ async fn merge_segments_to_dir( .join(INDEX_FILE_NAME) }) .collect::>(); - - let auxiliary_file = - lance_index::vector::distributed::index_merger::merge_partial_vector_auxiliary_files( - object_store, - &aux_paths, - final_dir, - progress.clone(), - ) - .await?; + let auxiliary_file = lance_index::vector::distributed::index_merger::merge_partial_vector_auxiliary_files_with_row_filters( + object_store, + &aux_paths, + final_dir, + row_filters, + progress.clone(), + ) + .await?; let index_file = write_root_vector_index_from_auxiliary( object_store, final_dir, From 6b8f4a09cf5c695bfb6a033462f96fefe17c52a9 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Sat, 8 Aug 2026 12:58:49 +0000 Subject: [PATCH 5/7] fix(index): preserve safe vector merge paths --- .../src/vector/distributed/index_merger.rs | 81 +++++++-- rust/lance/src/index.rs | 98 ++++++---- rust/lance/src/index/append.rs | 170 ++++++++++++++++-- rust/lance/src/io/exec/knn.rs | 9 +- 4 files changed, 295 insertions(+), 63 deletions(-) diff --git a/rust/lance-index/src/vector/distributed/index_merger.rs b/rust/lance-index/src/vector/distributed/index_merger.rs index eb91ac46432..5cfcfba49ac 100755 --- a/rust/lance-index/src/vector/distributed/index_merger.rs +++ b/rust/lance-index/src/vector/distributed/index_merger.rs @@ -429,6 +429,35 @@ pub async fn write_partition_rows( Ok(()) } +/// Stream a partition range, retain its owned rows, and return the number written. +async fn write_filtered_partition_rows( + reader: &V2Reader, + w: &mut FileWriter, + range: Range, + row_filter: &OldIndexDataFilter, +) -> Result { + let mut stream = reader + .read_stream( + lance_io::ReadBatchParams::Range(range), + u32::MAX, + 4, + lance_encoding::decoder::FilterExpression::no_filter(), + ) + .await?; + let mut written_rows = 0usize; + while let Some(batch) = stream.next().await { + let batch = filter_batch_to_owned_rows(&batch?, row_filter)?; + if batch.num_rows() == 0 { + continue; + } + written_rows = written_rows.checked_add(batch.num_rows()).ok_or_else(|| { + Error::index("Filtered partition row count exceeds usize capacity".to_string()) + })?; + w.write_batch(&batch).await?; + } + Ok(written_rows) +} + /// Transpose the PQ code column for a batch and write it to the unified writer. /// /// This helper assumes `batch` contains a contiguous range of rows for a single @@ -1625,25 +1654,45 @@ async fn merge_partial_vector_auxiliary_files_inner( } } _ => { - let partition_window_size = *PARTITION_WINDOW_SIZE; - let prefetch_window_count = *PARTITION_PREFETCH_WINDOW_COUNT; - let mut shard_merge_reader = ShardMergeReader::new( - shard_infos, - nlist, - partition_window_size, - prefetch_window_count, - ); - while let Some((pid, batches)) = shard_merge_reader.next_partition().await? { - let partition_len = batches.iter().map(RecordBatch::num_rows).sum::(); + // FLAT, SQ, and their HNSW variants do not need whole-partition + // transforms. Stream one shard partition at a time so filtering + // never materializes a multi-partition window in memory. + for (pid, merged_length) in merged_lengths.iter_mut().enumerate() { + let mut partition_len = 0usize; + for shard in &shard_infos { + let source_len = shard.lengths[pid] as usize; + if source_len == 0 { + continue; + } + let offset = shard.partition_offsets[pid]; + let writer = v2w_opt.as_mut().ok_or_else(|| { + Error::index("Failed to initialize unified writer".to_string()) + })?; + let written = if let Some(row_filter) = shard.row_filter.as_deref() { + write_filtered_partition_rows( + shard.reader.as_ref(), + writer, + offset..offset + source_len, + row_filter, + ) + .await? + } else { + write_partition_rows( + shard.reader.as_ref(), + writer, + offset..offset + source_len, + ) + .await?; + source_len + }; + partition_len = partition_len.checked_add(written).ok_or_else(|| { + Error::index(format!("Merged partition {pid} exceeds usize row capacity")) + })?; + } if partition_len == 0 { continue; } - if let Some(w) = v2w_opt.as_mut() { - for batch in batches { - w.write_batch(&batch).await?; - } - } - merged_lengths[pid] = u32::try_from(partition_len).map_err(|_| { + *merged_length = u32::try_from(partition_len).map_err(|_| { Error::index(format!("Merged partition {pid} exceeds u32 row capacity")) })?; merged_rows = merged_rows.saturating_add(partition_len as u64); diff --git a/rust/lance/src/index.rs b/rust/lance/src/index.rs index 9733b996cf1..eb57c0e4d7a 100644 --- a/rust/lance/src/index.rs +++ b/rust/lance/src/index.rs @@ -65,7 +65,8 @@ use tracing::{info, instrument, warn}; use uuid::Uuid; use vector::details::{ derive_vector_index_type, infer_missing_vector_details, merged_physical_fragment_bitmap, - needs_vector_details_inference, vector_details_as_json, with_physical_fragment_bitmap, + needs_vector_details_inference, physical_fragment_bitmap, vector_details_as_json, + with_physical_fragment_bitmap, }; pub(crate) use vector::details::{vector_index_details, vector_index_details_default}; use vector::ivf::v2::{IVFIndex, IvfStateEntryBox}; @@ -136,17 +137,30 @@ fn project_index_through_fragment_reuse( index: &mut IndexMetadata, frag_reuse_index: &FragReuseIndex, ) -> Result<()> { + let invalidates_physical_coverage = physical_fragment_bitmap(index).is_some_and(|coverage| { + frag_reuse_index.details.versions.iter().any(|version| { + version.dataset_version >= index.dataset_version + && version.groups.iter().any(|group| { + group + .old_frags + .iter() + .any(|fragment| coverage.contains(fragment.id as u32)) + }) + }) + }); let Some(fragment_bitmap) = index.fragment_bitmap.as_mut() else { return Ok(()); }; frag_reuse_index.remap_fragment_bitmap(fragment_bitmap)?; - index.index_details = index - .index_details - .as_deref() - .cloned() - .map(|details| with_physical_fragment_bitmap(details, None)) - .transpose()? - .map(Arc::new); + if invalidates_physical_coverage { + index.index_details = index + .index_details + .as_deref() + .cloned() + .map(|details| with_physical_fragment_bitmap(details, None)) + .transpose()? + .map(Arc::new); + } Ok(()) } @@ -3233,25 +3247,28 @@ mod tests { use std::collections::{HashMap, HashSet}; #[test] - fn test_fragment_reuse_projection_invalidates_vector_physical_coverage() { - let physical_fragments = RoaringBitmap::from_iter([0_u32]); - let details = prost_types::Any::from_msg(&VectorIndexDetails { - compression: Some(Compression::Flat(FlatCompression {})), - ..Default::default() - }) - .unwrap(); - let details = with_physical_fragment_bitmap(details, Some(&physical_fragments)).unwrap(); - let mut index = IndexMetadata { - uuid: Uuid::new_v4(), - name: "vector_idx".to_string(), - fields: vec![0], - dataset_version: 1, - fragment_bitmap: Some(physical_fragments), - index_details: Some(Arc::new(details)), - index_version: 1, - created_at: None, - base_id: None, - files: None, + fn test_fragment_reuse_projection_invalidates_only_affected_vector_coverage() { + let vector_segment = |dataset_version, fragment_id| { + let physical_fragments = RoaringBitmap::from_iter([fragment_id]); + let details = prost_types::Any::from_msg(&VectorIndexDetails { + compression: Some(Compression::Flat(FlatCompression {})), + ..Default::default() + }) + .unwrap(); + let details = + with_physical_fragment_bitmap(details, Some(&physical_fragments)).unwrap(); + IndexMetadata { + uuid: Uuid::new_v4(), + name: "vector_idx".to_string(), + fields: vec![0], + dataset_version, + fragment_bitmap: Some(physical_fragments), + index_details: Some(Arc::new(details)), + index_version: 1, + created_at: None, + base_id: None, + files: None, + } }; let digest = |id| FragDigest { id, @@ -3273,11 +3290,18 @@ mod tests { }, ); - project_index_through_fragment_reuse(&mut index, &frag_reuse_index).unwrap(); + let mut affected_index = vector_segment(1, 0); + project_index_through_fragment_reuse(&mut affected_index, &frag_reuse_index).unwrap(); - assert_eq!(index.fragment_bitmap, Some(RoaringBitmap::from_iter([1]))); - assert_eq!(vector::details::physical_fragment_bitmap(&index), None); - let remapped_details = index + assert_eq!( + affected_index.fragment_bitmap, + Some(RoaringBitmap::from_iter([1])) + ); + assert_eq!( + vector::details::physical_fragment_bitmap(&affected_index), + None + ); + let remapped_details = affected_index .index_details .as_deref() .unwrap() @@ -3287,6 +3311,18 @@ mod tests { remapped_details.compression, Some(Compression::Flat(_)) )); + + let mut newer_disjoint_index = vector_segment(3, 10); + project_index_through_fragment_reuse(&mut newer_disjoint_index, &frag_reuse_index).unwrap(); + assert_eq!( + newer_disjoint_index.fragment_bitmap, + Some(RoaringBitmap::from_iter([10])) + ); + assert_eq!( + vector::details::physical_fragment_bitmap(&newer_disjoint_index), + Some(RoaringBitmap::from_iter([10])), + "an older disjoint reuse mapping must preserve exact provenance" + ); } async fn write_vector_segment_metadata( diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index db27cd1732a..f6f1a2a97d1 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -690,19 +690,29 @@ async fn build_fresh_vector_segment( .await } -fn segment_merge_requires_rebuild( - index: &IndexMetadata, - current_fragments: &RoaringBitmap, -) -> bool { - let Some(physical_fragments) = physical_fragment_bitmap(index) else { - return true; - }; +fn segment_merge_requires_rebuild(dataset: &Dataset, index: &IndexMetadata) -> bool { let owned_fragments = index - .effective_fragment_bitmap(current_fragments) + .effective_fragment_bitmap(&dataset.fragment_bitmap) .or_else(|| index.fragment_bitmap.clone()) .unwrap_or_default(); + + if dataset.manifest.uses_stable_row_ids() + && dataset.get_fragments().iter().any(|fragment| { + owned_fragments.contains(fragment.id() as u32) + && fragment.metadata().deletion_file.is_some() + }) + { + // Fragment coverage cannot distinguish deleted and replacement rows + // that share a stable row id. Rebuild this source from current rows so + // stale values cannot be copied into the merged segment. + return true; + } + + let Some(physical_fragments) = physical_fragment_bitmap(index) else { + return true; + }; let unowned_physical_fragments = &physical_fragments - &owned_fragments; - unowned_physical_fragments.intersection_len(current_fragments) > 0 + unowned_physical_fragments.intersection_len(&dataset.fragment_bitmap) > 0 } async fn scan_vector_fragments( @@ -1025,7 +1035,7 @@ pub async fn merge_indices_with_unindexed_frags<'a>( .effective_fragment_bitmap(&dataset.fragment_bitmap) .or_else(|| removed_segment.fragment_bitmap.clone()) .unwrap_or_default(); - if segment_merge_requires_rebuild(removed_segment, &dataset.fragment_bitmap) { + if segment_merge_requires_rebuild(dataset.as_ref(), removed_segment) { let segment = build_fresh_vector_segment( dataset.as_ref(), &selected_logical_index, @@ -1092,9 +1102,10 @@ pub async fn merge_indices_with_unindexed_frags<'a>( .collect(), )?; let selected_segments = &old_indices[merge_start..]; - if selected_segments.iter().any(|segment| { - segment_merge_requires_rebuild(segment, &dataset.fragment_bitmap) - }) { + if selected_segments + .iter() + .any(|segment| segment_merge_requires_rebuild(dataset.as_ref(), segment)) + { let mut fragment_bitmap = base_unindexed_bitmap.clone(); for segment in selected_segments { if let Some(effective) = @@ -2801,6 +2812,139 @@ mod tests { assert_eq!(results[0].num_rows(), 10); } + #[tokio::test] + async fn test_vector_merge_filters_stable_row_id_replacements() { + const DIMENSION: usize = 4; + + let test_dir = TempStrDir::default(); + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new( + "vector", + DataType::FixedSizeList( + Arc::new(Field::new("item", DataType::Float32, true)), + DIMENSION as i32, + ), + false, + ), + ])); + let initial_values = (0..40) + .flat_map(|row| [if row < 20 { 1.0 } else { 0.0 }; DIMENSION]) + .collect::>(); + let initial = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from_iter_values(0..40)), + Arc::new( + FixedSizeListArray::try_new_from_values( + arrow_array::Float32Array::from(initial_values), + DIMENSION as i32, + ) + .unwrap(), + ), + ], + ) + .unwrap(); + let mut dataset = Dataset::write( + RecordBatchIterator::new([Ok(initial)], schema.clone()), + test_dir.as_str(), + Some(WriteParams { + enable_stable_row_ids: true, + ..Default::default() + }), + ) + .await + .unwrap(); + dataset + .create_index( + &["vector"], + IndexType::Vector, + Some("vector_idx".to_string()), + &VectorIndexParams::ivf_flat(1, MetricType::L2), + true, + ) + .await + .unwrap(); + + let replacements = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from_iter_values(20..40)), + Arc::new( + FixedSizeListArray::try_new_from_values( + arrow_array::Float32Array::from(vec![10.0; 20 * DIMENSION]), + DIMENSION as i32, + ) + .unwrap(), + ), + ], + ) + .unwrap(); + let merge_job = MergeInsertBuilder::try_new(Arc::new(dataset), vec!["id".to_string()]) + .unwrap() + .when_matched(WhenMatched::UpdateAll) + .try_build() + .unwrap(); + let (dataset, stats) = merge_job + .execute(reader_to_stream(Box::new(RecordBatchIterator::new( + [Ok(replacements)], + schema, + )))) + .await + .unwrap(); + assert_eq!(stats.num_updated_rows, 20); + + let mut dataset = dataset.as_ref().clone(); + dataset + .optimize_indices(&OptimizeOptions::append()) + .await + .unwrap(); + assert_eq!( + dataset + .load_indices_by_name("vector_idx") + .await + .unwrap() + .len(), + 2 + ); + dataset + .optimize_indices(&OptimizeOptions::merge(2)) + .await + .unwrap(); + + let logical_index = dataset + .open_logical_vector_index("vector", "vector_idx") + .await + .unwrap(); + assert_eq!(logical_index.num_segments(), 1); + assert_eq!( + logical_index + .num_rows_per_segment() + .into_iter() + .map(|(_, rows)| rows) + .sum::(), + 40, + "the merged index must contain one current copy of every stable row id" + ); + + let query = arrow_array::Float32Array::from(vec![0.0; DIMENSION]); + let result = dataset + .scan() + .project(&["id"]) + .unwrap() + .nearest("vector", &query, 5) + .unwrap() + .nprobes(1) + .try_into_batch() + .await + .unwrap(); + let ids = result["id"].as_primitive::(); + assert!( + ids.values().iter().all(|id| *id < 20), + "stale pre-update vectors must not survive the optimize merge: {ids:?}" + ); + } + #[tokio::test] async fn test_merge_indices_with_unindexed_frags_vector_subset() { const DIM: usize = 64; diff --git a/rust/lance/src/io/exec/knn.rs b/rust/lance/src/io/exec/knn.rs index a40196f706a..b031e725fa3 100644 --- a/rust/lance/src/io/exec/knn.rs +++ b/rust/lance/src/io/exec/knn.rs @@ -2972,13 +2972,16 @@ mod tests { dataset.append(second, None).await.unwrap(); let appended_fragments = dataset.fragment_bitmap.as_ref() - &first_fragments; let dataset = Arc::new(dataset); + let configured_details = crate::index::vector::details::vector_index_details( + &VectorIndexParams::ivf_flat(1, MetricType::L2), + ); let old_details = crate::index::vector::details::with_physical_fragment_bitmap( - crate::index::vector_index_details_default(), + configured_details.clone(), Some(&first_fragments), ) .unwrap(); let new_details = crate::index::vector::details::with_physical_fragment_bitmap( - crate::index::vector_index_details_default(), + configured_details.clone(), Some(&appended_fragments), ) .unwrap(); @@ -3014,7 +3017,7 @@ mod tests { let merged_physical_fragments = &first_fragments | &appended_fragments; let merged_details = crate::index::vector::details::with_physical_fragment_bitmap( - crate::index::vector_index_details_default(), + configured_details, Some(&merged_physical_fragments), ) .unwrap(); From 1c7ff252f6b992bf6c7c12409ba0dd572e5ce3a3 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Sat, 8 Aug 2026 13:34:08 +0000 Subject: [PATCH 6/7] fix(index): preserve clean stable vector merge paths --- rust/lance/src/index/append.rs | 57 ++++++++++++++++++++++++++++++++-- 1 file changed, 54 insertions(+), 3 deletions(-) diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index f6f1a2a97d1..fdf780a016e 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -699,12 +699,19 @@ fn segment_merge_requires_rebuild(dataset: &Dataset, index: &IndexMetadata) -> b if dataset.manifest.uses_stable_row_ids() && dataset.get_fragments().iter().any(|fragment| { owned_fragments.contains(fragment.id() as u32) - && fragment.metadata().deletion_file.is_some() + && fragment + .metadata() + .deletion_file + .as_ref() + .is_some_and(|deletion_file| { + deletion_file.read_version >= index.dataset_version + }) }) { // Fragment coverage cannot distinguish deleted and replacement rows - // that share a stable row id. Rebuild this source from current rows so - // stale values cannot be copied into the merged segment. + // that share a stable row id. A deletion older than the segment was + // already applied when it was built; rebuild only for later deletions + // so stale values cannot be copied into the merged segment. return true; } @@ -2945,6 +2952,50 @@ mod tests { ); } + #[tokio::test] + async fn test_stable_row_id_segment_built_after_deletion_keeps_merge_fast_path() { + let mut dataset = lance_datagen::gen_batch() + .col("id", array::step::()) + .col("vector", array::rand_vec::(Dimension::from(4))) + .into_dataset_with_params( + "memory://", + FragmentCount(1), + FragmentRowCount(40), + Some(WriteParams { + enable_stable_row_ids: true, + max_rows_per_file: 40, + ..Default::default() + }), + ) + .await + .unwrap(); + dataset.delete("id < 20").await.unwrap(); + dataset + .create_index( + &["vector"], + IndexType::Vector, + Some("vector_idx".to_string()), + &VectorIndexParams::ivf_flat(1, MetricType::L2), + true, + ) + .await + .unwrap(); + + let segment = dataset + .load_indices_by_name("vector_idx") + .await + .unwrap() + .pop() + .unwrap(); + let fragments = dataset.get_fragments(); + let deletion_file = fragments[0].metadata().deletion_file.as_ref().unwrap(); + assert!(deletion_file.read_version < segment.dataset_version); + assert!( + !segment_merge_requires_rebuild(&dataset, &segment), + "a segment built after a deletion must retain the auxiliary merge fast path" + ); + } + #[tokio::test] async fn test_merge_indices_with_unindexed_frags_vector_subset() { const DIM: usize = 64; From ca23eed82ae4d2bfb8a9d83b7d9e327d206734af Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Sat, 8 Aug 2026 15:04:39 +0000 Subject: [PATCH 7/7] fix(index): preserve stable vector merge topology --- rust/lance/src/index/append.rs | 19 ------------------- 1 file changed, 19 deletions(-) diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 4dfcbe18627..0f36add0170 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -696,25 +696,6 @@ fn segment_merge_requires_rebuild(dataset: &Dataset, index: &IndexMetadata) -> b .or_else(|| index.fragment_bitmap.clone()) .unwrap_or_default(); - if dataset.manifest.uses_stable_row_ids() - && dataset.get_fragments().iter().any(|fragment| { - owned_fragments.contains(fragment.id() as u32) - && fragment - .metadata() - .deletion_file - .as_ref() - .is_some_and(|deletion_file| { - deletion_file.read_version >= index.dataset_version - }) - }) - { - // Fragment coverage cannot distinguish deleted and replacement rows - // that share a stable row id. A deletion older than the segment was - // already applied when it was built; rebuild only for later deletions - // so stale values cannot be copied into the merged segment. - return true; - } - let Some(physical_fragments) = physical_fragment_bitmap(index) else { return true; };