From a4317f99a70bf3677372c1b7162ecfcf67b49294 Mon Sep 17 00:00:00 2001 From: Will Jones Date: Wed, 22 Jul 2026 16:32:38 -0700 Subject: [PATCH] fix(index): keep data overlays masked through OptimizeIndices Overlay index masking gates on `overlay.committed_version > segment.dataset_version`. The merge path of `OptimizeIndices` carries an old segment's stale pre-overlay entries without re-reading overlaid values, but stamped the merged segment with the current `dataset_version`. That flipped the gate off and un-masked the overlay: a stale index entry resurfaced and the overlaid value was dropped. This affected every index type (BTREE, BITMAP, FTS, IVF vector). Treat an overlay-stale fragment (touched by a field-relevant overlay committed after the segment was built) like a compaction-retired one during merge, mirroring `prune_overlay_stale_fields_from_indices`. The new `segment_coverage_split` drops such fragments from the merged segment's coverage; scalar types (BTree/bitmap/label_list/FTS) additionally filter their old entries out via `OldIndexDataFilter`. Vector needs only the coverage exclusion because its ANN prefilter is already restricted to coverage (`DatasetPreFilter::new`), whereas the scalar `MaterializeIndexExec` prefilter is not. Pruned fragments fall to the flat path and read current (overlay-merged) values. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../tests/dataset_overlay_index_masking.rs | 229 ++++++++++++++++++ rust/lance/src/index/append.rs | 115 ++++++--- 2 files changed, 315 insertions(+), 29 deletions(-) diff --git a/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs b/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs index cfabdc13266..65a80376917 100644 --- a/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs +++ b/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs @@ -1055,3 +1055,232 @@ async fn bench_index_query_overlay_overhead() { println!("{num_vec_overlays:>12} {ann_ms:>10.1}"); } } + +async fn append_age_fragment(dataset: &mut Dataset, ids: std::ops::Range) { + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("id", DataType::Int32, true), + ArrowField::new("age", DataType::Int32, true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from_iter_values(ids.clone())), + Arc::new(Int32Array::from_iter_values(ids.map(|v| v * 10))), + ], + ) + .unwrap(); + dataset + .append( + RecordBatchIterator::new(vec![Ok(batch)], schema.clone()), + None, + ) + .await + .unwrap(); +} + +// `OptimizeIndices` merges an index's delta segments without re-reading data overlays, but used +// to stamp the merged segment with the current `dataset_version`. That flipped the overlay mask's +// version gate (`overlay.committed_version > segment.dataset_version`) off, un-masking the stale +// pre-overlay entries the merge had carried over. The merge now drops overlay-stale fragments from +// the new segment's coverage (and, for scalar types, filters their old entries out via +// `OldIndexDataFilter`) so those fragments fall to the flat path against current values. The +// following tests reproduce the un-masking for each index type and assert it stays masked. + +/// Scalar (BTree, Bitmap): a range query over the indexed column after an overlay + optimize must +/// still drop the stale value and surface the overlaid one. +#[rstest] +#[case::btree(IndexType::BTree)] +#[case::bitmap(IndexType::Bitmap)] +#[tokio::test] +async fn test_optimize_preserves_scalar_overlay_masking(#[case] index_type: IndexType) { + use crate::index::DatasetIndexExt; + use lance_index::optimize::OptimizeOptions; + use lance_index::scalar::BuiltinIndexType; + + let params = match index_type { + IndexType::Bitmap => ScalarIndexParams::for_builtin(BuiltinIndexType::Bitmap), + _ => ScalarIndexParams::default(), + }; + let mut dataset = create_base_dataset().await; + dataset + .create_index(&["age"], index_type, None, ¶ms, true) + .await + .unwrap(); + + // Overlay fragment 0, offset 1 (id=1): age 10 -> 999, committed after the index. + let mut dataset = commit_overlay( + dataset, + "age_opt", + 0, + &[1], + OverlayCoverage::dense(RoaringBitmap::from_iter([1])), + vec![i32_array([Some(999)])], + ) + .await; + + // Masking works before optimize. + assert_eq!(ids_matching(&dataset, "age = 10").await, Vec::::new()); + assert_eq!(ids_matching(&dataset, "age = 999").await, vec![1]); + + // Append an unindexed fragment so the merge does real work, then merge all deltas. + append_age_fragment(&mut dataset, 12..18).await; + dataset + .optimize_indices(&OptimizeOptions::merge(10)) + .await + .unwrap(); + + // Still masked: the stale age=10 entry stays dropped and the overlaid age=999 value stays + // visible (fragment 0 is dropped from index coverage and re-evaluated on the flat path). + assert_eq!( + ids_matching(&dataset, "age = 10").await, + Vec::::new(), + "stale index entry age=10 for id=1 resurfaced after optimize" + ); + assert_eq!( + ids_matching(&dataset, "age = 999").await, + vec![1], + "overlaid value age=999 dropped after optimize" + ); + // A row untouched by the overlay is unaffected. + assert_eq!(ids_matching(&dataset, "age = 20").await, vec![2]); +} + +/// FTS: after an overlay replaces a row's text and the index is optimized, searching for the old +/// terms must not return the stale row, and the new terms must find it. +#[tokio::test] +async fn test_optimize_preserves_fts_overlay_masking() { + use crate::index::DatasetIndexExt; + use lance_index::optimize::OptimizeOptions; + + let mut dataset = create_text_dataset().await; + build_text_fts_index(&mut dataset).await; + + // fragment 0, offset 1 (id=1): "apple banana" -> "cherry mango". + let mut dataset = commit_overlay( + dataset, + "text_opt", + 0, + &[1], + OverlayCoverage::dense(RoaringBitmap::from_iter([1])), + vec![Arc::new(StringArray::from(vec![Some("cherry mango")]))], + ) + .await; + + // Masking works before optimize: id=1 no longer matches "banana"/"apple". + assert_eq!(fts_ids_matching(&dataset, "banana").await, vec![3]); + assert_eq!(fts_ids_matching(&dataset, "apple").await, vec![0]); + + // Append an unindexed fragment of new text, then merge all deltas. + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("id", DataType::Int32, true), + ArrowField::new("text", DataType::Utf8, true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from_iter_values(12..18)), + Arc::new(StringArray::from(vec![ + "kiwi", "melon", "date", "guava", "papaya", "lychee", + ])), + ], + ) + .unwrap(); + dataset + .append( + RecordBatchIterator::new(vec![Ok(batch)], schema.clone()), + None, + ) + .await + .unwrap(); + dataset + .optimize_indices(&OptimizeOptions::merge(10)) + .await + .unwrap(); + + // Still masked: id=1's stale "apple"/"banana" postings stay dropped. + assert_eq!( + fts_ids_matching(&dataset, "banana").await, + vec![3], + "stale FTS posting for id=1 (banana) resurfaced after optimize" + ); + assert_eq!( + fts_ids_matching(&dataset, "apple").await, + vec![0], + "stale FTS posting for id=1 (apple) resurfaced after optimize" + ); + // The overlaid terms are found via the flat path. + assert!(fts_ids_matching(&dataset, "cherry").await.contains(&1)); + assert!(fts_ids_matching(&dataset, "mango").await.contains(&1)); +} + +/// Vector (IVF): after an overlay moves a row's vector and the index is optimized, the ANN must +/// not resurface the stale vector, and the moved-onto-query row is found by flat re-scoring. +#[tokio::test] +async fn test_optimize_preserves_vector_overlay_masking() { + use crate::index::DatasetIndexExt; + use lance_index::optimize::OptimizeOptions; + + // Overlay on fragment 1 moves id=35 away from the query and id=40 onto it. + let mut dataset = create_vector_overlay_dataset(false).await; + + // Masking works before optimize. + let before = vector_query_ids(&dataset, 3, false).await; + assert!( + !before.contains(&35), + "pre-optimize id=35 should be dropped: {before:?}" + ); + assert!( + before.contains(&40), + "pre-optimize id=40 should be found: {before:?}" + ); + + // Append an unindexed fragment of far vectors, then merge all deltas. + let far_vecs: Vec> = (0..32) + .map(|i| { + let mut v = vec![0.0_f32; VEC_DIM as usize]; + v[1] = (i + 200) as f32; + v + }) + .collect(); + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new("id", DataType::Int32, true), + ArrowField::new( + "vec", + DataType::FixedSizeList( + Arc::new(ArrowField::new("item", DataType::Float32, true)), + VEC_DIM, + ), + true, + ), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from_iter_values(64..96)), + fsl(far_vecs, VEC_DIM), + ], + ) + .unwrap(); + dataset + .append( + RecordBatchIterator::new(vec![Ok(batch)], schema.clone()), + None, + ) + .await + .unwrap(); + dataset + .optimize_indices(&OptimizeOptions::merge(10)) + .await + .unwrap(); + + // Still masked: id=35's stale vector stays dropped and id=40 is still found via re-scoring. + let after = vector_query_ids(&dataset, 3, false).await; + assert!( + !after.contains(&35), + "stale index vector for id=35 resurfaced after optimize: {after:?}" + ); + assert!( + after.contains(&40), + "overlaid vector for id=40 dropped after optimize: {after:?}" + ); +} diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 9dfc61894b2..8e3cbcdd596 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -26,11 +26,14 @@ use lance_table::format::{Fragment, IndexMetadata}; use roaring::RoaringBitmap; use uuid::Uuid; +use std::collections::HashMap; + use super::DatasetIndexInternalExt; use super::vector::LogicalVectorIndex; use super::vector::ivf::{optimize_vector_indices, select_segment_for_single_rebalance}; use crate::dataset::Dataset; use crate::dataset::index::LanceIndexStoreExt; +use crate::dataset::overlay::{collect_overlay_stale_frags, overlaid_fragments}; use crate::dataset::rowids::load_row_id_sequences; use crate::index::scalar::{IndexDetails, fetch_index_details, load_training_data}; use crate::index::vector_index_details_default; @@ -133,24 +136,55 @@ pub async fn build_old_data_filter( } } -/// Split the stored fragment coverage of `segments` into fragments still live in -/// `dataset` (`effective`) and fragments that compaction or deletion has already -/// retired (`deleted`). +/// One segment's live coverage split into fragments the merge keeps on the indexed path +/// (`effective`) and fragments whose old index entries must be dropped (`retired`). +/// +/// `retired` unions two causes: fragments no longer in the dataset (compaction/deletion), and +/// fragments carrying a data overlay this merge will *not* incorporate -- one committed after the +/// segment was built (`committed_version > segment.dataset_version`) that touches a field the +/// segment indexes. A merge carries a segment's stored entries without re-reading overlaid values, +/// so such a fragment's indexed values are stale; dropping it from coverage and filtering its old +/// entries out routes the whole fragment to the flat path, which reads current (overlay-merged) +/// values. Mirrors `prune_overlay_stale_fields_from_indices`, which does the same for overlays a +/// compaction has baked into rewritten fragments. +fn segment_coverage_split( + dataset: &Dataset, + segment: &IndexMetadata, + overlaid_frags: &HashMap, +) -> Result<(RoaringBitmap, RoaringBitmap)> { + let mut effective = segment + .effective_fragment_bitmap(&dataset.fragment_bitmap) + .unwrap_or_default(); + let mut retired = segment + .deleted_fragment_bitmap(&dataset.fragment_bitmap) + .unwrap_or_default(); + let mut stale = RoaringBitmap::new(); + collect_overlay_stale_frags(segment, overlaid_frags, &mut stale)?; + // Only still-live fragments matter for coverage; a stale fragment already retired above is a + // no-op in both directions. + stale &= &effective; + effective -= &stale; + retired |= stale; + Ok((effective, retired)) +} + +/// Split the fragment coverage of `segments` into fragments the merge keeps on the indexed path +/// (`effective`) and fragments whose entries must be dropped (`retired`): those retired by +/// compaction/deletion, plus those carrying an overlay the merge won't incorporate. See +/// [`segment_coverage_split`]. pub fn split_segment_coverage<'a>( dataset: &Dataset, segments: impl IntoIterator, -) -> (RoaringBitmap, RoaringBitmap) { +) -> Result<(RoaringBitmap, RoaringBitmap)> { + let overlaid = overlaid_fragments(&dataset.manifest.fragments); let mut effective = RoaringBitmap::new(); - let mut deleted = RoaringBitmap::new(); + let mut retired = RoaringBitmap::new(); for segment in segments { - if let Some(eff) = segment.effective_fragment_bitmap(&dataset.fragment_bitmap) { - effective |= eff; - } - if let Some(del) = segment.deleted_fragment_bitmap(&dataset.fragment_bitmap) { - deleted |= del; - } + let (eff, ret) = segment_coverage_split(dataset, segment, &overlaid)?; + effective |= eff; + retired |= ret; } - (effective, deleted) + Ok((effective, retired)) } /// Build one [`OldIndexDataFilter`] per segment, each derived from that segment's @@ -160,6 +194,7 @@ pub async fn build_per_segment_filters( dataset: &Dataset, segments: &[&IndexMetadata], ) -> Result<(RoaringBitmap, Vec>)> { + let overlaid = overlaid_fragments(&dataset.manifest.fragments); let mut effective_union = RoaringBitmap::new(); let mut filters = Vec::with_capacity(segments.len()); for segment in segments { @@ -169,14 +204,9 @@ pub async fn build_per_segment_filters( segment.uuid ))); } - let effective = segment - .effective_fragment_bitmap(&dataset.fragment_bitmap) - .unwrap_or_default(); - let deleted = segment - .deleted_fragment_bitmap(&dataset.fragment_bitmap) - .unwrap_or_default(); + let (effective, retired) = segment_coverage_split(dataset, segment, &overlaid)?; effective_union |= &effective; - filters.push(build_old_data_filter(dataset, &effective, &deleted).await?); + filters.push(build_old_data_filter(dataset, &effective, &retired).await?); } Ok((effective_union, filters)) } @@ -406,9 +436,10 @@ async fn merge_scalar_indices<'a>( .await?; let update_criteria = reference_index.update_criteria(); - // Effective = bitmap ∩ live fragments; deleted = bitmap \ live fragments. - let (effective_old_frags, deleted_old_frags) = - split_segment_coverage(dataset.as_ref(), selected_old_indices.iter().copied()); + // Effective = fragments the merge keeps indexed; retired = fragments whose old entries must + // be dropped (retired by compaction/deletion, or stale under an unincorporated overlay). + let (effective_old_frags, retired_old_frags) = + split_segment_coverage(dataset.as_ref(), selected_old_indices.iter().copied())?; let mut frag_bitmap = base_unindexed_bitmap.clone(); frag_bitmap |= &effective_old_frags; @@ -498,7 +529,7 @@ async fn merge_scalar_indices<'a>( let old_data_filter = build_old_data_filter( dataset.as_ref(), &effective_old_frags, - &deleted_old_frags, + &retired_old_frags, ) .await?; reference_index @@ -866,15 +897,19 @@ pub async fn merge_indices_with_unindexed_frags<'a>( ) .await?; + let overlaid = overlaid_fragments(&dataset.manifest.fragments); let mut frag_bitmap = base_unindexed_bitmap; let mut effective_old_frags = RoaringBitmap::new(); + let mut retired_old_frags = RoaringBitmap::new(); let mut selected_indices = Vec::with_capacity(selected_old_indices.len()); for idx in &selected_old_indices { - if let Some(effective) = idx.effective_fragment_bitmap(&dataset.fragment_bitmap) - { - frag_bitmap |= &effective; - effective_old_frags |= &effective; - } + // Drop fragments carrying an unincorporated overlay from coverage (see + // `segment_coverage_split`); they fall to the flat FTS path with current values. + let (effective, retired) = + segment_coverage_split(dataset.as_ref(), idx, &overlaid)?; + frag_bitmap |= &effective; + effective_old_frags |= &effective; + retired_old_frags |= retired; let scalar_index = dataset .open_scalar_index(&field_path, &idx.uuid, &NoOpMetricsCollector) .await?; @@ -900,7 +935,7 @@ pub async fn merge_indices_with_unindexed_frags<'a>( } else { Some(OldIndexDataFilter::Fragments { to_keep: effective_old_frags, - to_remove: RoaringBitmap::new(), + to_remove: retired_old_frags, }) }; @@ -961,6 +996,28 @@ pub async fn merge_indices_with_unindexed_frags<'a>( } }?; + // A vector merge carries the old segments' quantized codes without re-reading overlaid values, + // and vector indices have no `OldIndexDataFilter` to drop stale postings. Instead, drop + // overlay-stale fragments from the new segment's coverage: the ANN prefilter is an allow-list + // restricted to coverage (see `DatasetPreFilter::new`), so their stale codes are blocked, and + // `knn_combined` re-scores those fragments on the flat path against current (overlay-merged) + // values. The scalar path already excludes these fragments (and filters their entries) inside + // `merge_scalar_indices`. + let new_fragment_bitmap = if first_is_vector_index { + let overlaid = overlaid_fragments(&dataset.manifest.fragments); + if overlaid.is_empty() { + new_fragment_bitmap + } else { + let mut stale = RoaringBitmap::new(); + for &segment in &removed_indices { + collect_overlay_stale_frags(segment, &overlaid, &mut stale)?; + } + &new_fragment_bitmap - &stale + } + } else { + new_fragment_bitmap + }; + Ok(Some(IndexMergeResults { new_uuid, removed_indices,