From 9fce59f5c74ab8e32839c9500ec89baf7bf23aff Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Wed, 3 Jun 2026 22:55:36 +0800 Subject: [PATCH 1/2] feat: add segmented BTree index merge_segments support --- java/src/test/java/org/lance/DatasetTest.java | 24 +- rust/lance-index/src/scalar/btree.rs | 104 ++- rust/lance/src/dataset.rs | 3 +- rust/lance/src/index.rs | 21 +- rust/lance/src/index/append.rs | 678 +++++++++++++++--- rust/lance/src/index/create.rs | 248 +++++-- rust/lance/src/index/scalar.rs | 1 + rust/lance/src/index/scalar/btree.rs | 193 +++++ 8 files changed, 1075 insertions(+), 197 deletions(-) create mode 100644 rust/lance/src/index/scalar/btree.rs diff --git a/java/src/test/java/org/lance/DatasetTest.java b/java/src/test/java/org/lance/DatasetTest.java index 3ea6a0812e1..45466a0367c 100644 --- a/java/src/test/java/org/lance/DatasetTest.java +++ b/java/src/test/java/org/lance/DatasetTest.java @@ -1993,18 +1993,20 @@ void testOptimizingIndices(@TempDir Path tempDir) throws Exception { OptimizeOptions options = OptimizeOptions.builder().numIndicesToMerge(0).build(); dsAppended.optimizeIndices(options); - List afterIndexes = dsAppended.getIndexes(); - Index idIndexAfter = - afterIndexes.stream() + List idIndexes = + dsAppended.getIndexes().stream() .filter(idx -> "id_idx".equals(idx.name())) - .findFirst() - .orElse(null); - assertNotNull(idIndexAfter); - List afterFragments = idIndexAfter.fragments().orElse(Collections.emptyList()); - - assertTrue(afterFragments.contains(0)); - assertTrue(afterFragments.contains(1)); - assertEquals(2, afterFragments.size()); + .collect(Collectors.toList()); + assertEquals( + 2, + idIndexes.size(), + "append-only optimize must add a delta segment instead of merging"); + + Set coveredFragments = + idIndexes.stream() + .flatMap(idx -> idx.fragments().orElse(Collections.emptyList()).stream()) + .collect(Collectors.toSet()); + assertEquals(new HashSet<>(Arrays.asList(0, 1)), coveredFragments); } } } diff --git a/rust/lance-index/src/scalar/btree.rs b/rust/lance-index/src/scalar/btree.rs index ba6d3dd142d..a650e8db244 100644 --- a/rust/lance-index/src/scalar/btree.rs +++ b/rust/lance-index/src/scalar/btree.rs @@ -68,7 +68,7 @@ use tracing::{info, instrument}; mod flat; -const BTREE_LOOKUP_NAME: &str = "page_lookup.lance"; +pub const BTREE_LOOKUP_NAME: &str = "page_lookup.lance"; const BTREE_PAGES_NAME: &str = "page_data.lance"; pub const DEFAULT_BTREE_BATCH_SIZE: u64 = 4096; const BATCH_SIZE_META_KEY: &str = "batch_size"; @@ -1489,7 +1489,7 @@ impl BTreeIndex { } /// Create a stream of all the data in the index, in the same format used to train the index - async fn into_data_stream(self) -> Result { + async fn data_stream(&self) -> Result { let lazy_reader = LazyIndexReader::new(self.store.clone(), self.ranges_to_files.clone()); let reader = lazy_reader.get().await?; let new_schema = Arc::new(self.train_schema()); @@ -1512,25 +1512,51 @@ impl BTreeIndex { ))) } - async fn combine_old_new( - self, + /// Merge N source BTree segments plus an additional `new_data` stream into + /// a single BTree under `dest_store`, without re-reading the dataset. + pub async fn merge_segments( + segments: &[Arc], new_data: SendableRecordBatchStream, - chunk_size: u64, + dest_store: &dyn IndexStore, old_data_filter: Option, - ) -> Result { - let value_column_index = new_data.schema().index_of(VALUE_COLUMN_NAME)?; - - let new_input = Arc::new(OneShotExec::new(new_data)); - let old_stream = self.into_data_stream().await?; - let old_stream = match old_data_filter { - Some(filter) => filter_row_ids(old_stream, filter), - None => old_stream, + ) -> Result { + let Some(first) = segments.first() else { + return Err(Error::invalid_input( + "cannot merge BTree index without at least one source segment".to_string(), + )); }; - let old_input = Arc::new(OneShotExec::new(old_stream)); - debug_assert_eq!( - old_input.schema().flattened_fields().len(), - new_input.schema().flattened_fields().len() - ); + + for segment in segments.iter().skip(1) { + if segment.data_type != first.data_type { + return Err(Error::index(format!( + "cannot merge BTree segments with different value types ({:?} vs {:?})", + first.data_type, segment.data_type + ))); + } + } + + let new_schema = new_data.schema(); + let value_column_index = new_schema.index_of(VALUE_COLUMN_NAME)?; + let new_value_type = new_schema.field(value_column_index).data_type(); + if new_value_type != &first.data_type { + return Err(Error::invalid_input(format!( + "BTree merge: new_data value column type {:?} does not match \ + segment value type {:?}", + new_value_type, first.data_type + ))); + } + + let mut inputs: Vec> = Vec::with_capacity(segments.len() + 1); + for segment in segments { + let stream = segment.data_stream().await?; + let stream = match old_data_filter.clone() { + Some(filter) => filter_row_ids(stream, filter), + None => stream, + }; + let exec = Arc::new(OneShotExec::new(stream)); + inputs.push(exec); + } + inputs.push(Arc::new(OneShotExec::new(new_data))); let sort_expr = PhysicalSortExpr { expr: Arc::new(Column::new(VALUE_COLUMN_NAME, value_column_index)), @@ -1539,11 +1565,10 @@ impl BTreeIndex { nulls_first: true, }, }; - // The UnionExec creates multiple partitions but the SortPreservingMergeExec merges - // them back into a single partition. - let all_data = UnionExec::try_new(vec![old_input, new_input])?; - let ordered = Arc::new(SortPreservingMergeExec::new([sort_expr].into(), all_data)); - + // UnionExec yields multiple partitions; SortPreservingMergeExec merges + // them back into a single partition while preserving value-ordering. + let unioned = UnionExec::try_new(inputs)?; + let ordered = Arc::new(SortPreservingMergeExec::new([sort_expr].into(), unioned)); let unchunked = execute_plan( ordered, LanceExecutionOptions { @@ -1551,7 +1576,16 @@ impl BTreeIndex { ..Default::default() }, )?; - Ok(chunk_concat_stream(unchunked, chunk_size as usize)) + let merged_stream = chunk_concat_stream(unchunked, first.batch_size as usize); + + train_btree_index(merged_stream, dest_store, first.batch_size, None, None).await?; + + Ok(CreatedIndex { + index_details: prost_types::Any::from_msg(&pbold::BTreeIndexDetails::default()) + .unwrap(), + index_version: BTREE_INDEX_VERSION, + files: Some(dest_store.list_files_with_sizes().await?), + }) } } @@ -1896,19 +1930,15 @@ impl ScalarIndex for BTreeIndex { dest_store: &dyn IndexStore, old_data_filter: Option, ) -> Result { - // Merge the existing index data with the new data and then retrain the index on the merged stream - let merged_data_source = self - .clone() - .combine_old_new(new_data, self.batch_size, old_data_filter) - .await?; - train_btree_index(merged_data_source, dest_store, self.batch_size, None, None).await?; - - Ok(CreatedIndex { - index_details: prost_types::Any::from_msg(&pbold::BTreeIndexDetails::default()) - .unwrap(), - index_version: BTREE_INDEX_VERSION, - files: Some(dest_store.list_files_with_sizes().await?), - }) + // Updating is the single-segment case of a segment merge: union this + // index's data with `new_data`, re-sort on value, and retrain. + Self::merge_segments( + &[Arc::new(self.clone())], + new_data, + dest_store, + old_data_filter, + ) + .await } fn update_criteria(&self) -> UpdateCriteria { diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 52e21bd7cba..74112332b61 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -3054,7 +3054,8 @@ impl Dataset { IndexType::BTree => { Err(Error::invalid_input( "BTree distributed indexing no longer supports merge_index_metadata; \ - build segments, and commit with commit_existing_index_segments(...)" + build segments, optionally merge groups with merge_existing_index_segments(...), \ + and commit with commit_existing_index_segments(...)" .to_string(), )) } diff --git a/rust/lance/src/index.rs b/rust/lance/src/index.rs index eaa3dc6119d..8c84d6b157c 100644 --- a/rust/lance/src/index.rs +++ b/rust/lance/src/index.rs @@ -47,6 +47,7 @@ use lance_index::{INDEX_FILE_NAME, Index, IndexType, PrewarmOptions, pb, vector: use lance_index::{ IndexCriteria, is_system_index, metrics::{MetricsCollector, NoOpMetricsCollector}, + scalar::btree::BTREE_LOOKUP_NAME, }; use lance_io::scheduler::{ScanScheduler, SchedulerConfig}; use lance_io::traits::Reader; @@ -250,6 +251,19 @@ fn segment_has_inverted_details(segment: &IndexMetadata) -> bool { .is_some_and(|details| details.type_url.ends_with("InvertedIndexDetails")) } +/// Detect BTree segments, preserving a legacy pre-details fallback. +fn segment_has_btree_details(segment: &IndexMetadata) -> bool { + segment.index_details.as_ref().map_or_else( + || { + segment + .files + .as_ref() + .is_some_and(|files| files.iter().any(|file| file.path == BTREE_LOOKUP_NAME)) + }, + |details| details.type_url.ends_with("BTreeIndexDetails"), + ) +} + // Cache keys for different index types #[derive(Debug, Clone)] pub(crate) struct LegacyVectorIndexCacheKey<'a> { @@ -1069,7 +1083,8 @@ impl DatasetIndexExt for Dataset { } let all_vector = source_segments.iter().all(segment_has_vector_details); let all_inverted = source_segments.iter().all(segment_has_inverted_details); - if !all_vector && !all_inverted { + let all_btree = source_segments.iter().all(segment_has_btree_details); + if !all_vector && !all_inverted && !all_btree { return Err(Error::invalid_input( "merge_existing_index_segments requires all segments to have the same supported index type" .to_string(), @@ -1083,8 +1098,10 @@ impl DatasetIndexExt for Dataset { source_segments, ) .await? - } else { + } else if all_inverted { crate::index::scalar::inverted::merge_segments(self, source_segments).await? + } else { + crate::index::scalar::btree::merge_segments(self, source_segments).await? }; merged_segment.dataset_version = self.manifest.version; merged_segment.fields = vec![field_id]; diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 4398928d3e2..257babb535a 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -74,6 +74,200 @@ async fn build_stable_row_id_filter( Ok(::union_all(&row_id_map_refs)) } +/// Build the [`OldIndexDataFilter`] that must be applied to existing index +/// rows when their owning fragments have been pruned by compaction or +/// deletions. +pub async fn build_old_data_filter( + dataset: &Dataset, + effective_old_frags: &RoaringBitmap, + deleted_old_frags: &RoaringBitmap, +) -> Result> { + if dataset.manifest.uses_stable_row_ids() { + let valid_old_row_ids = build_stable_row_id_filter(dataset, effective_old_frags).await?; + Ok(Some(OldIndexDataFilter::RowIds(valid_old_row_ids))) + } else { + Ok(Some(OldIndexDataFilter::Fragments { + to_keep: effective_old_frags.clone(), + to_remove: deleted_old_frags.clone(), + })) + } +} + +async fn load_unindexed_training_data( + dataset: &Dataset, + field_path: &str, + update_criteria: &lance_index::scalar::UpdateCriteria, + unindexed: &[Fragment], +) -> Result { + let fragments = if update_criteria.requires_old_data { + None + } else { + Some(unindexed.to_vec()) + }; + load_training_data( + dataset, + field_path, + &update_criteria.data_criteria, + fragments, + true, + None, + ) + .await +} + +#[allow(clippy::too_many_arguments)] +async fn merge_scalar_indices<'a>( + dataset: Arc, + old_indices: &[&'a IndexMetadata], + unindexed: &[Fragment], + options: &OptimizeOptions, + index_type: IndexType, + field_path: &str, + column_name: &str, + base_unindexed_bitmap: RoaringBitmap, +) -> Result, RoaringBitmap, CreatedIndex)>> { + if old_indices.is_empty() { + return Err(Error::index( + "merge_scalar_indices: no previous index found".to_string(), + )); + } + + let num_to_merge = match index_type { + IndexType::BTree => options + .num_indices_to_merge + .unwrap_or(1) + .min(old_indices.len()), + _ => { + if old_indices.len() > 1 { + return Err(Error::not_supported(format!( + "Cannot optimize multi-segment {:?} scalar index \ + ({} segments): this index type lacks a k-way \ + segment merge primitive. Drop and recreate the \ + index to consolidate.", + index_type, + old_indices.len(), + ))); + } + 1 + } + }; + + // No new data + ≤1 old selected = rewriting one segment to itself. + if unindexed.is_empty() && num_to_merge <= 1 { + return Ok(None); + } + + let selected_old_indices = &old_indices[old_indices.len() - num_to_merge..]; + + // For the delta case (`selected` empty) the reference is purely + // for reading params; fall back to the last old index then. + let reference_idx = selected_old_indices + .first() + .copied() + .unwrap_or(old_indices[old_indices.len() - 1]); + let reference_index = dataset + .open_scalar_index( + field_path, + &reference_idx.uuid.to_string(), + &NoOpMetricsCollector, + ) + .await?; + + // Effective = bitmap ∩ live fragments; deleted = bitmap \ live fragments. + let mut effective_old_frags = RoaringBitmap::new(); + let mut deleted_old_frags = RoaringBitmap::new(); + for idx in selected_old_indices { + if let Some(effective) = idx.effective_fragment_bitmap(&dataset.fragment_bitmap) { + effective_old_frags |= effective; + } + if let Some(deleted) = idx.deleted_fragment_bitmap(&dataset.fragment_bitmap) { + deleted_old_frags |= deleted; + } + } + + let mut frag_bitmap = base_unindexed_bitmap.clone(); + frag_bitmap |= &effective_old_frags; + let new_uuid = Uuid::new_v4(); + + let created_index = if selected_old_indices.is_empty() { + // Delta: build from new_data only. Old segments stay in manifest. + let params = reference_index.derive_index_params()?; + let unindexed_fragment_ids: Vec = base_unindexed_bitmap.iter().collect(); + super::scalar::build_scalar_index( + dataset.as_ref(), + column_name, + &new_uuid.to_string(), + ¶ms, + true, + Some(unindexed_fragment_ids), + None, + Arc::new(NoopIndexBuildProgress), + ) + .await? + } else { + let update_criteria = reference_index.update_criteria(); + let new_data_stream = + load_unindexed_training_data(dataset.as_ref(), field_path, &update_criteria, unindexed) + .await?; + + if effective_old_frags.is_empty() { + // Stale: every selected segment's coverage was retired. + // Rebuild from new_data; `selected` is returned as + // `removed_indices` so the manifest drops them. + let params = reference_index.derive_index_params()?; + super::scalar::build_scalar_index( + dataset.as_ref(), + column_name, + &new_uuid.to_string(), + ¶ms, + true, + None, + Some(new_data_stream), + Arc::new(NoopIndexBuildProgress), + ) + .await? + } else { + let new_store = LanceIndexStore::from_dataset_for_new(&dataset, &new_uuid.to_string())?; + let old_data_filter = + build_old_data_filter(dataset.as_ref(), &effective_old_frags, &deleted_old_frags) + .await?; + + match index_type { + IndexType::BTree => { + crate::index::scalar::btree::open_and_merge_segments( + dataset.as_ref(), + field_path, + selected_old_indices, + new_data_stream, + &new_store, + old_data_filter, + ) + .await? + } + _ => { + debug_assert_eq!( + selected_old_indices.len(), + 1, + "num_to_merge clamp invariant violated: index types \ + without a k-way merge primitive must reach \ + update() with exactly one selected segment" + ); + reference_index + .update(new_data_stream, &new_store, old_data_filter) + .await? + } + } + } + }; + + Ok(Some(( + new_uuid, + selected_old_indices.to_vec(), + frag_bitmap, + created_index, + ))) +} + async fn metadata_is_vector_index(dataset: &Dataset, index: &IndexMetadata) -> Result { if let Some(files) = &index.files { return Ok(files.iter().any(|file| file.path == INDEX_FILE_NAME)); @@ -339,7 +533,6 @@ pub async fn merge_indices_with_unindexed_frags<'a>( )) } } else { - let mut frag_bitmap = base_unindexed_bitmap.clone(); let mut indices = Vec::with_capacity(old_indices.len()); for idx in old_indices { match dataset @@ -515,105 +708,21 @@ pub async fn merge_indices_with_unindexed_frags<'a>( )) } it if it.is_scalar() => { - let num_to_merge = options - .num_indices_to_merge - .unwrap_or(1) - .min(old_indices.len()); - if unindexed.is_empty() && num_to_merge <= 1 { - return Ok(None); - } - - // Use effective bitmap (intersected with existing dataset fragments) - // to avoid carrying stale data from pruned indices. - let effective_old_frags: RoaringBitmap = old_indices - .iter() - .filter_map(|idx| idx.effective_fragment_bitmap(&dataset.fragment_bitmap)) - .fold(RoaringBitmap::new(), |mut acc, b| { - acc |= &b; - acc - }); - let deleted_old_frags: RoaringBitmap = old_indices - .iter() - .filter_map(|idx| idx.deleted_fragment_bitmap(&dataset.fragment_bitmap)) - .fold(RoaringBitmap::new(), |mut acc, b| { - acc |= &b; - acc - }); - frag_bitmap |= &effective_old_frags; - - let index = dataset - .open_scalar_index( - &field_path, - &old_indices[0].uuid.to_string(), - &NoOpMetricsCollector, - ) - .await?; - - let update_criteria = index.update_criteria(); - - let fragments = if update_criteria.requires_old_data { - None - } else { - Some(unindexed.to_vec()) - }; - let new_data_stream = load_training_data( - dataset.as_ref(), + let Some(result) = merge_scalar_indices( + dataset.clone(), + old_indices, + unindexed, + options, + it, &field_path, - &update_criteria.data_criteria, - fragments, - true, - None, + column.name.as_str(), + base_unindexed_bitmap, ) - .await?; - - let new_uuid = Uuid::new_v4(); - - let created_index = if effective_old_frags.is_empty() { - // Old data is fully stale (bitmap pruned to empty). Rebuild - // from scratch instead of merging stale entries. - let params = index.derive_index_params()?; - super::scalar::build_scalar_index( - dataset.as_ref(), - column.name.as_str(), - &new_uuid.to_string(), - ¶ms, - true, - None, - Some(new_data_stream), - Arc::new(NoopIndexBuildProgress), - ) - .await? - } else { - let new_store = - LanceIndexStore::from_dataset_for_new(&dataset, &new_uuid.to_string())?; - let old_data_filter = if dataset.manifest.uses_stable_row_ids() { - // Stable row IDs are opaque IDs, so fragment-bit filtering on - // (row_id >> 32) is invalid. Build an exact allow-list from retained - // fragments' row-id sequences and use precise filtering. - let valid_old_row_ids = - build_stable_row_id_filter(dataset.as_ref(), &effective_old_frags) - .await?; - Some(OldIndexDataFilter::RowIds(valid_old_row_ids)) - } else { - // Address-style row IDs encode fragment_id in high 32 bits. - // Fragment bitmap filtering is valid and cheaper in this mode. - Some(OldIndexDataFilter::Fragments { - to_keep: effective_old_frags, - to_remove: deleted_old_frags, - }) - }; - index - .update(new_data_stream, &new_store, old_data_filter) - .await? + .await? + else { + return Ok(None); }; - - // TODO: don't hard-code index version - Ok(( - new_uuid, - vec![old_indices[old_indices.len() - 1]], - frag_bitmap, - created_index, - )) + Ok(result) } _ => Err(Error::index(format!( "Append index: invalid index type: {:?}", @@ -662,7 +771,7 @@ mod tests { use crate::dataset::builder::DatasetBuilder; use crate::dataset::optimize::compact_files; - use crate::dataset::{MergeInsertBuilder, WhenMatched, WhenNotMatched, WriteParams}; + use crate::dataset::{MergeInsertBuilder, WhenMatched, WhenNotMatched, WriteMode, WriteParams}; use crate::index::vector::VectorIndexParams; use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; @@ -1219,6 +1328,379 @@ mod tests { ); } + #[tokio::test] + async fn test_optimize_btree_multi_segment_optimize_default() { + async fn query_id_count(dataset: &Dataset, id: &str) -> usize { + dataset + .scan() + .filter(&format!("id = '{}'", id)) + .unwrap() + .project(&["id"]) + .unwrap() + .try_into_batch() + .await + .unwrap() + .num_rows() + } + + let test_dir = TempStrDir::default(); + let test_uri = test_dir.as_str(); + + let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Utf8, false)])); + let make_batch = |start: i32, end: i32| { + let ids = StringArray::from_iter_values((start..end).map(|i| format!("song-{i}"))); + RecordBatch::try_new(schema.clone(), vec![Arc::new(ids)]).unwrap() + }; + + // Three fragments of 64 rows each; each commits as its own BTree + // segment so optimize sees a multi-segment scalar logical index. + let reader = RecordBatchIterator::new( + vec![ + Ok(make_batch(0, 64)), + Ok(make_batch(64, 128)), + Ok(make_batch(128, 192)), + ], + schema.clone(), + ); + let mut dataset = Dataset::write( + reader, + test_uri, + Some(WriteParams { + max_rows_per_file: 64, + ..Default::default() + }), + ) + .await + .unwrap(); + + let params = ScalarIndexParams::for_builtin(lance_index::scalar::BuiltinIndexType::BTree); + let fragments = dataset.get_fragments(); + assert_eq!(fragments.len(), 3); + + let mut staged_segments = Vec::new(); + for fragment in &fragments { + let segment = crate::index::create::CreateIndexBuilder::new( + &mut dataset, + &["id"], + IndexType::BTree, + ¶ms, + ) + .name("id_idx".into()) + .fragments(vec![fragment.id() as u32]) + .execute_uncommitted() + .await + .unwrap(); + staged_segments.push(segment); + } + dataset + .commit_existing_index_segments("id_idx", "id", staged_segments) + .await + .unwrap(); + assert_eq!( + dataset.load_indices_by_name("id_idx").await.unwrap().len(), + 3 + ); + + let appended = RecordBatchIterator::new(vec![Ok(make_batch(192, 256))], schema.clone()); + let mut dataset = Dataset::write( + appended, + test_uri, + Some(WriteParams { + max_rows_per_file: 64, + mode: WriteMode::Append, + ..Default::default() + }), + ) + .await + .unwrap(); + assert_eq!(dataset.get_fragments().len(), 4); + + dataset + .optimize_indices(&OptimizeOptions::default()) + .await + .unwrap(); + + // Reload from disk to ensure we're reading committed manifest state. + let dataset = DatasetBuilder::from_uri(test_uri).load().await.unwrap(); + + // Each of these IDs lives in a distinct old segment / fragment. + // song-10 lives in fragment 0, song-80 in fragment 1, song-160 in + // fragment 2, and song-200 in the appended fragment. After optimize + // every row must still be reachable through the logical index, + // regardless of which segment absorbed the new data. + for id in ["song-10", "song-80", "song-160", "song-200"] { + assert_eq!( + query_id_count(&dataset, id).await, + 1, + "expected exactly one row for {id} after multi-segment optimize" + ); + } + + // `OptimizeOptions::default()` (= num_indices_to_merge: None) merges + // the newest segment with the unindexed fragment, like the + // inverted/vector default. The three old segments minus the merged one + // plus the new delta means three segments remain, and together they + // must still cover every dataset fragment without overlap. + let segments_after = dataset.load_indices_by_name("id_idx").await.unwrap(); + assert_eq!( + segments_after.len(), + 3, + "default optimize must merge one delta, not all segments, got {segments_after:?}" + ); + let mut covered = RoaringBitmap::new(); + for segment in &segments_after { + let bitmap = segment + .fragment_bitmap + .as_ref() + .expect("each segment should carry fragment coverage"); + assert!( + covered.is_disjoint(bitmap), + "post-optimize segments must not overlap, got {segments_after:?}" + ); + covered |= bitmap; + } + let mut expected = RoaringBitmap::new(); + for frag in dataset.get_fragments() { + expected.insert(frag.id() as u32); + } + assert_eq!( + covered, expected, + "post-optimize segments should cover every dataset fragment" + ); + } + + #[tokio::test] + async fn test_optimize_btree_optimize_append() { + async fn query_id_count(dataset: &Dataset, id: &str) -> usize { + dataset + .scan() + .filter(&format!("id = '{}'", id)) + .unwrap() + .project(&["id"]) + .unwrap() + .try_into_batch() + .await + .unwrap() + .num_rows() + } + + let test_dir = TempStrDir::default(); + let test_uri = test_dir.as_str(); + + let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Utf8, false)])); + let make_batch = |start: i32, end: i32| { + let ids = StringArray::from_iter_values((start..end).map(|i| format!("song-{i}"))); + RecordBatch::try_new(schema.clone(), vec![Arc::new(ids)]).unwrap() + }; + + // Start with two fragments + two committed BTree segments. + let reader = RecordBatchIterator::new( + vec![Ok(make_batch(0, 64)), Ok(make_batch(64, 128))], + schema.clone(), + ); + let mut dataset = Dataset::write( + reader, + test_uri, + Some(WriteParams { + max_rows_per_file: 64, + ..Default::default() + }), + ) + .await + .unwrap(); + + let params = ScalarIndexParams::for_builtin(lance_index::scalar::BuiltinIndexType::BTree); + let original_segment_uuids: Vec<_> = { + let mut staged = Vec::new(); + for fragment in dataset.get_fragments() { + let segment = crate::index::create::CreateIndexBuilder::new( + &mut dataset, + &["id"], + IndexType::BTree, + ¶ms, + ) + .name("id_idx".into()) + .fragments(vec![fragment.id() as u32]) + .execute_uncommitted() + .await + .unwrap(); + staged.push(segment); + } + let uuids = staged.iter().map(|s| s.uuid).collect::>(); + dataset + .commit_existing_index_segments("id_idx", "id", staged) + .await + .unwrap(); + uuids + }; + assert_eq!(original_segment_uuids.len(), 2); + + // Append a third fragment, leave it unindexed, then run append-mode optimize. + let appended = RecordBatchIterator::new(vec![Ok(make_batch(128, 192))], schema.clone()); + let mut dataset = Dataset::write( + appended, + test_uri, + Some(WriteParams { + max_rows_per_file: 64, + mode: WriteMode::Append, + ..Default::default() + }), + ) + .await + .unwrap(); + + dataset + .optimize_indices(&OptimizeOptions::append()) + .await + .unwrap(); + + // Read fresh from disk to make sure we're inspecting committed state. + let dataset = DatasetBuilder::from_uri(test_uri).load().await.unwrap(); + + // append() must preserve every original old segment unchanged and add + // exactly one new segment covering only the newly appended fragments. + let committed = dataset.load_indices_by_name("id_idx").await.unwrap(); + let committed_uuids: std::collections::HashSet<_> = + committed.iter().map(|idx| idx.uuid).collect(); + for original in &original_segment_uuids { + assert!( + committed_uuids.contains(original), + "append() must not remove pre-existing segment {original}, \ + but the committed UUIDs are {committed_uuids:?}" + ); + } + assert_eq!( + committed.len(), + original_segment_uuids.len() + 1, + "append() should add exactly one new delta segment, got {committed:?}" + ); + let new_segment = committed + .iter() + .find(|idx| !original_segment_uuids.contains(&idx.uuid)) + .expect("append() must add a new delta segment"); + let new_segment_frags: Vec<_> = new_segment + .fragment_bitmap + .as_ref() + .unwrap() + .iter() + .collect(); + // The appended fragment should be the only one covered by the new delta; + // old segments retain their own coverage. + assert_eq!(new_segment_frags.len(), 1); + + // Sanity check: queries across all fragments still return their rows. + for id in ["song-10", "song-100", "song-160"] { + assert_eq!(query_id_count(&dataset, id).await, 1, "missing row {id}"); + } + } + + #[tokio::test] + async fn test_optimize_bitmap_index_append() { + let test_dir = TempStrDir::default(); + let test_uri = test_dir.as_str(); + + let schema = Arc::new(Schema::new(vec![Field::new( + "category", + DataType::Utf8, + false, + )])); + let make_batch = |labels: &[&str]| { + let arr = StringArray::from_iter_values(labels.iter().copied()); + RecordBatch::try_new(schema.clone(), vec![Arc::new(arr)]).unwrap() + }; + + // One fragment + one Bitmap segment. + let reader = + RecordBatchIterator::new(vec![Ok(make_batch(&["a", "b", "a", "c"]))], schema.clone()); + let mut dataset = Dataset::write( + reader, + test_uri, + Some(WriteParams { + max_rows_per_file: 4, + ..Default::default() + }), + ) + .await + .unwrap(); + + let params = ScalarIndexParams::for_builtin(lance_index::scalar::BuiltinIndexType::Bitmap); + dataset + .create_index( + &["category"], + IndexType::Bitmap, + Some("cat_idx".into()), + ¶ms, + true, + ) + .await + .unwrap(); + let original_uuid = { + let committed = dataset.load_indices_by_name("cat_idx").await.unwrap(); + assert_eq!(committed.len(), 1); + committed[0].uuid + }; + + // Append a second fragment, leave it unindexed, then optimize with + // `append()` (= num_indices_to_merge: Some(0)). + let appended = + RecordBatchIterator::new(vec![Ok(make_batch(&["b", "d", "d", "a"]))], schema.clone()); + let mut dataset = Dataset::write( + appended, + test_uri, + Some(WriteParams { + max_rows_per_file: 4, + mode: WriteMode::Append, + ..Default::default() + }), + ) + .await + .unwrap(); + + dataset + .optimize_indices(&OptimizeOptions::append()) + .await + .unwrap(); + let dataset = DatasetBuilder::from_uri(test_uri).load().await.unwrap(); + + // Bitmap remains a single-segment logical index. The original UUID + // is replaced (because update() produces a new UUID), but we never + // grow to 2 segments. + let committed = dataset.load_indices_by_name("cat_idx").await.unwrap(); + assert_eq!( + committed.len(), + 1, + "Bitmap optimize append() must silently merge instead of creating a 2nd segment, \ + got {committed:?}" + ); + // The merged segment should cover both fragments (old + new). + let frags: std::collections::BTreeSet = committed[0] + .fragment_bitmap + .as_ref() + .expect("merged Bitmap should carry fragment coverage") + .iter() + .collect(); + assert_eq!(frags, [0u32, 1].into_iter().collect()); + // Sanity: the original UUID was replaced by the merge. + assert_ne!( + committed[0].uuid, original_uuid, + "update() should produce a new UUID, but committed segment still has the original" + ); + + // Data correctness: a value that lives only in the appended fragment + // is queryable through the index. + let rows = dataset + .scan() + .filter("category = 'd'") + .unwrap() + .project(&["category"]) + .unwrap() + .try_into_batch() + .await + .unwrap() + .num_rows(); + assert_eq!(rows, 2, "value 'd' lives in appended fragment"); + } + #[tokio::test] async fn test_optimize_btree_keeps_rows_with_stable_row_ids_after_compaction() { async fn query_id_count(dataset: &Dataset, id: &str) -> usize { diff --git a/rust/lance/src/index/create.rs b/rust/lance/src/index/create.rs index af7ea7ce19c..e21a4643f12 100644 --- a/rust/lance/src/index/create.rs +++ b/rust/lance/src/index/create.rs @@ -480,7 +480,7 @@ impl<'a> CreateIndexBuilder<'a> { } else { vec![] }; - let transaction = if uses_segment_commit_path(self.index_type, &new_idx.name, self.params) { + let transaction = if uses_segment_commit_path(self.index_type, self.params) { let field_id = *new_idx.fields.first().ok_or_else(|| { Error::internal(format!( "Index '{}' is missing field ids after build", @@ -534,6 +534,13 @@ impl<'a> CreateIndexBuilder<'a> { } } +fn is_btree_scalar_params(params: &dyn IndexParams) -> bool { + params + .as_any() + .downcast_ref::() + .is_some_and(|p| p.index_type.eq_ignore_ascii_case("btree")) +} + /// Validate that a user-supplied `index_uuid` is permitted for this build. fn ensure_index_uuid_allowed( index_type: IndexType, @@ -560,26 +567,35 @@ fn ensure_index_uuid_allowed( Ok(()) } -fn uses_segment_commit_path( - index_type: IndexType, - index_name: &str, - params: &dyn IndexParams, -) -> bool { - if index_name != LANCE_VECTOR_INDEX { - return false; +fn uses_segment_commit_path(index_type: IndexType, params: &dyn IndexParams) -> bool { + let params_family = params.index_name(); + + if params_family == LANCE_VECTOR_INDEX + && matches!( + index_type, + IndexType::Vector + | IndexType::IvfPq + | IndexType::IvfSq + | IndexType::IvfFlat + | IndexType::IvfRq + | IndexType::IvfHnswFlat + | IndexType::IvfHnswPq + | IndexType::IvfHnswSq + ) + && params.as_any().is::() + { + return true; } - matches!( - index_type, - IndexType::Vector - | IndexType::IvfPq - | IndexType::IvfSq - | IndexType::IvfFlat - | IndexType::IvfRq - | IndexType::IvfHnswFlat - | IndexType::IvfHnswPq - | IndexType::IvfHnswSq - ) && params.as_any().is::() + if params_family == LANCE_SCALAR_INDEX { + match index_type { + IndexType::BTree => return true, + IndexType::Scalar if is_btree_scalar_params(params) => return true, + _ => {} + } + } + + false } impl<'a> IntoFuture for CreateIndexBuilder<'a> { @@ -1933,6 +1949,143 @@ mod tests { assert_eq!(results.num_rows(), 20); } + #[tokio::test] + async fn test_btree_merge_existing_index_segments() { + use datafusion::common::ScalarValue; + use lance_index::scalar::{SargableQuery, SearchResult}; + use std::ops::Bound; + + // Open `segment` and count rows whose `id` falls in `[lo, hi)`. + async fn count_in_range( + dataset: &Dataset, + segment: &IndexMetadata, + lo: i32, + hi: i32, + ) -> usize { + let field_path = dataset.schema().field_path(segment.fields[0]).unwrap(); + let index = crate::index::scalar::open_scalar_index( + dataset, + &field_path, + segment, + &NoOpMetricsCollector, + ) + .await + .unwrap(); + let query = SargableQuery::Range( + Bound::Included(ScalarValue::Int32(Some(lo))), + Bound::Excluded(ScalarValue::Int32(Some(hi))), + ); + match index.search(&query, &NoOpMetricsCollector).await.unwrap() { + SearchResult::Exact(row_addrs) => { + row_addrs.true_rows().row_addrs().unwrap().count() + } + other => panic!("expected exact result, got {other:?}"), + } + } + + let tmpdir = TempStrDir::default(); + let dataset_uri = format!("file://{}", tmpdir.as_str()); + + // 128 rows across two 64-row fragments. Stable row ids so the + // retired-fragment filter below exercises the exact row-id allow-list. + let reader = gen_batch() + .col("id", lance_datagen::array::step::()) + .into_reader_rows( + lance_datagen::RowCount::from(64), + lance_datagen::BatchCount::from(2), + ); + let mut dataset = Dataset::write( + reader, + &dataset_uri, + Some(WriteParams { + max_rows_per_file: 64, + mode: WriteMode::Overwrite, + enable_stable_row_ids: true, + ..Default::default() + }), + ) + .await + .unwrap(); + + // One staged BTree segment per fragment, committed as a multi-segment + // logical index. + let params = ScalarIndexParams::for_builtin(lance_index::scalar::BuiltinIndexType::BTree); + let mut staged = Vec::new(); + for fragment in dataset.get_fragments() { + staged.push( + CreateIndexBuilder::new(&mut dataset, &["id"], IndexType::BTree, ¶ms) + .name("id_btree".to_string()) + .fragments(vec![fragment.id() as u32]) + .execute_uncommitted() + .await + .unwrap(), + ); + } + dataset + .commit_existing_index_segments("id_btree", "id", staged) + .await + .unwrap(); + + // Phase 1 — healthy merge: the two per-fragment segments consolidate + // into a single canonical segment covering both fragments, and a range + // spanning both (ids 50..100) returns every matching row. + let merged = dataset + .merge_existing_index_segments(dataset.load_indices_by_name("id_btree").await.unwrap()) + .await + .unwrap(); + assert_eq!( + merged.fragment_bitmap.as_ref().unwrap(), + &roaring::RoaringBitmap::from_iter([0u32, 1]) + ); + assert!( + merged + .index_details + .as_ref() + .unwrap() + .type_url + .ends_with("BTreeIndexDetails") + ); + assert_eq!(count_in_range(&dataset, &merged, 50, 100).await, 50); + + // Phase 2 — retire fragment 0: delete >10% of its rows so compaction + // rewrites only frag 0 (frag 1 has no deletions and is at target size). + // The committed per-fragment segment now claims a fragment the dataset + // no longer has. + dataset.delete("id < 16").await.unwrap(); + crate::dataset::optimize::compact_files( + &mut dataset, + crate::dataset::optimize::CompactionOptions { + target_rows_per_fragment: 64, + ..Default::default() + }, + None, + ) + .await + .unwrap(); + let live_frags: roaring::RoaringBitmap = dataset + .get_fragments() + .iter() + .map(|f| f.id() as u32) + .collect(); + assert!(!live_frags.contains(0), "compaction should retire frag 0"); + + // Filtered merge: coverage drops the retired fragment but keeps the + // live one, and the merged page data does not leak the retired row ids + // (ids < 16 lived only in frag 0, so the range now returns nothing). + let merged = dataset + .merge_existing_index_segments(dataset.load_indices_by_name("id_btree").await.unwrap()) + .await + .unwrap(); + let coverage = merged.fragment_bitmap.as_ref().unwrap(); + assert!(!coverage.contains(0), "must drop retired frag 0"); + assert!(coverage.contains(1), "must keep live frag 1"); + assert_eq!( + count_in_range(&dataset, &merged, 0, 16).await, + 0, + "must filter retired-fragment row ids" + ); + } + #[tokio::test] async fn test_commit_existing_index_supports_local_hnsw_segments() { let tmpdir = TempStrDir::default(); @@ -2194,39 +2347,38 @@ mod tests { // Load indices after optimization let indices_after = dataset.load_indices().await.unwrap(); - // There should be 3 indices: - // 1. one scalar index with name "id_idx", and the bitmap is [0,1] - // 2. one delta vector index with name "vector_idx", and the bitmap is [0] - // 3. one delta vector index with name "vector_idx", and the bitmap is [1] - assert_eq!(indices_after.len(), 3, "{:?}", indices_after); - let id_idx = indices_after + // After unifying scalar optimize, `OptimizeOptions::append()` honors + // `Some(0)` for BTree the same way it does for vector: keep the old + // segment, add a delta for the unindexed fragment. So we now expect: + // 1. id_idx old segment, bitmap [0] + // 2. id_idx delta segment, bitmap [1] + // 3. vector_idx old segment, bitmap [0] + // 4. vector_idx delta segment, bitmap [1] + // Previously BTree silently merged into 1 segment because legacy + // scalar ignored `num_indices_to_merge`. + assert_eq!(indices_after.len(), 4, "{:?}", indices_after); + let id_indices = indices_after .iter() - .find(|idx| idx.name == "id_idx") - .unwrap(); + .filter(|idx| idx.name == "id_idx") + .collect::>(); let vector_indices = indices_after .iter() .filter(|idx| idx.name == "vector_idx") .collect::>(); - assert!( - id_idx - .fragment_bitmap - .as_ref() - .unwrap() - .contains_range(0..2) - && id_idx.fragment_bitmap.as_ref().unwrap().len() == 2 - ); - assert_eq!(vector_indices.len(), 2); - assert!( - vector_indices - .iter() - .any(|idx| idx.fragment_bitmap.as_ref().unwrap().contains(0) - && idx.fragment_bitmap.as_ref().unwrap().len() == 1) - ); - assert!( - vector_indices - .iter() - .any(|idx| idx.fragment_bitmap.as_ref().unwrap().contains(1) - && idx.fragment_bitmap.as_ref().unwrap().len() == 1) - ); + for indices in [&id_indices, &vector_indices] { + assert_eq!(indices.len(), 2); + assert!( + indices + .iter() + .any(|idx| idx.fragment_bitmap.as_ref().unwrap().contains(0) + && idx.fragment_bitmap.as_ref().unwrap().len() == 1) + ); + assert!( + indices + .iter() + .any(|idx| idx.fragment_bitmap.as_ref().unwrap().contains(1) + && idx.fragment_bitmap.as_ref().unwrap().len() == 1) + ); + } } } diff --git a/rust/lance/src/index/scalar.rs b/rust/lance/src/index/scalar.rs index 18c218ef4f7..7133b512f8b 100644 --- a/rust/lance/src/index/scalar.rs +++ b/rust/lance/src/index/scalar.rs @@ -4,6 +4,7 @@ //! Utilities for integrating scalar indices with datasets //! +pub(crate) mod btree; pub(crate) mod inverted; pub use inverted::{load_segment_details, load_segments}; diff --git a/rust/lance/src/index/scalar/btree.rs b/rust/lance/src/index/scalar/btree.rs new file mode 100644 index 00000000000..34534f6811b --- /dev/null +++ b/rust/lance/src/index/scalar/btree.rs @@ -0,0 +1,193 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +#![allow(clippy::redundant_pub_crate)] + +//! BTree-specific helpers for the segmented index workflow. +use std::sync::Arc; + +use arrow_schema::{Field as ArrowField, Schema as ArrowSchema}; +use datafusion::execution::SendableRecordBatchStream; +use datafusion::physical_plan::stream::RecordBatchStreamAdapter; +use lance_core::ROW_ID; +use lance_index::metrics::NoOpMetricsCollector; +use lance_index::pbold::BTreeIndexDetails; +use lance_index::scalar::btree::BTreeIndex; +use lance_index::scalar::lance_format::LanceIndexStore; +use lance_index::scalar::registry::VALUE_COLUMN_NAME; +use lance_index::scalar::{CreatedIndex, OldIndexDataFilter}; +use lance_table::format::IndexMetadata; +use roaring::RoaringBitmap; +use uuid::Uuid; + +use crate::{Dataset, Error, Result, dataset::index::LanceIndexStoreExt}; + +/// Build a row-empty `new_data` stream for the BTree merge API. +fn empty_btree_update_stream( + dataset: &Dataset, + field_id: i32, +) -> Result { + let field = dataset.schema().field_by_id(field_id).ok_or_else(|| { + Error::invalid_input(format!( + "merge_existing_index_segments: field id {} does not exist", + field_id + )) + })?; + let schema = Arc::new(ArrowSchema::new(vec![ + ArrowField::new(VALUE_COLUMN_NAME, field.data_type(), true), + ArrowField::new(ROW_ID, arrow_schema::DataType::UInt64, false), + ])); + Ok(Box::pin(RecordBatchStreamAdapter::new( + schema, + futures::stream::empty(), + ))) +} + +fn ensure_btree_details(segment: &IndexMetadata) -> Result<()> { + if let Some(details) = segment.index_details.as_ref() + && !details.type_url.ends_with("BTreeIndexDetails") + { + return Err(Error::invalid_input(format!( + "Segment '{}' is not a BTree segment (details type_url = '{}')", + segment.uuid, details.type_url + ))); + } + Ok(()) +} + +/// Open the given BTree `segments` and k-way merge their already-sorted page +/// data, together with `new_data`, into a single canonical BTree written to +/// `new_store`. +pub(crate) async fn open_and_merge_segments( + dataset: &Dataset, + field_path: &str, + segments: &[&IndexMetadata], + new_data: SendableRecordBatchStream, + new_store: &LanceIndexStore, + old_data_filter: Option, +) -> Result { + let mut source_indices = Vec::with_capacity(segments.len()); + for &segment in segments { + let scalar_index = + super::open_scalar_index(dataset, field_path, segment, &NoOpMetricsCollector).await?; + let btree = scalar_index + .as_any() + .downcast_ref::() + .ok_or_else(|| { + Error::index(format!( + "BTree merge: expected BTree segment {}, got {:?}", + segment.uuid, + scalar_index.index_type() + )) + })?; + source_indices.push(Arc::new(btree.clone())); + } + BTreeIndex::merge_segments(&source_indices, new_data, new_store, old_data_filter).await +} + +/// Merge one caller-defined group of source BTree segments into a single +/// physical segment. +pub(crate) async fn merge_segments( + dataset: &Dataset, + segments: Vec, +) -> Result { + if segments.is_empty() { + return Err(Error::index("No segment metadata was provided".to_string())); + } + + for segment in &segments { + ensure_btree_details(segment)?; + } + + // All source segments must belong to the same column. + let reference_fields = segments[0].fields.as_slice(); + for segment in segments.iter().skip(1) { + if segment.fields.as_slice() != reference_fields { + return Err(Error::invalid_input(format!( + "BTree merge_segments: segment {} has fields {:?}, expected {:?}", + segment.uuid, segment.fields, reference_fields, + ))); + } + } + + let field_id = *segments[0].fields.first().ok_or_else(|| { + Error::invalid_input(format!( + "CreateIndex: segment {} is missing field ids", + segments[0].uuid + )) + })?; + let field_path = dataset.schema().field_path(field_id)?; + + // Intersect each segment's stored bitmap with the dataset's current + // fragments so we don't claim coverage on IDs that compaction or pruning + // has already retired. + let dataset_fragments = dataset.fragment_bitmap.as_ref(); + let mut effective_old_frags = RoaringBitmap::new(); + let mut deleted_old_frags = RoaringBitmap::new(); + for segment in &segments { + if segment.fragment_bitmap.is_none() { + return Err(Error::invalid_input(format!( + "CreateIndex: segment {} is missing fragment coverage", + segment.uuid + ))); + } + if let Some(effective) = segment.effective_fragment_bitmap(dataset_fragments) { + effective_old_frags |= effective; + } + if let Some(deleted) = segment.deleted_fragment_bitmap(dataset_fragments) { + deleted_old_frags |= deleted; + } + } + + let fragment_bitmap = effective_old_frags.clone(); + let old_data_filter = crate::index::append::build_old_data_filter( + dataset, + &effective_old_frags, + &deleted_old_frags, + ) + .await?; + + let output_uuid = Uuid::new_v4(); + let new_store = LanceIndexStore::from_dataset_for_new(dataset, &output_uuid.to_string())?; + // Pure segment consolidation: no dataset scan, so `new_data` is an empty + // stream and the merge is driven entirely by the source page data. + let empty_new_data = empty_btree_update_stream(dataset, field_id)?; + let segment_refs: Vec<&IndexMetadata> = segments.iter().collect(); + let created_index = open_and_merge_segments( + dataset, + &field_path, + &segment_refs, + empty_new_data, + &new_store, + old_data_filter, + ) + .await?; + + if !created_index + .index_details + .type_url + .ends_with("BTreeIndexDetails") + { + return Err(Error::internal(format!( + "merge_existing_index_segments: BTree merge produced unexpected details type_url '{}'", + created_index.index_details.type_url + ))); + } + debug_assert_eq!( + created_index.index_details, + prost_types::Any::from_msg(&BTreeIndexDetails::default()).unwrap(), + ); + + Ok(IndexMetadata { + uuid: output_uuid, + name: segments[0].name.clone(), + fields: vec![field_id], + dataset_version: dataset.manifest.version, + fragment_bitmap: Some(fragment_bitmap), + index_details: Some(Arc::new(created_index.index_details)), + index_version: created_index.index_version as i32, + created_at: Some(chrono::Utc::now()), + base_id: None, + files: created_index.files, + }) +} From 0ec3fdaeff7cc4c923968d6d2e7bd7aa3614462c Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Thu, 4 Jun 2026 22:07:40 +0800 Subject: [PATCH 2/2] code review --- rust/lance/src/index/append.rs | 189 +++++++++++++++++---------------- 1 file changed, 99 insertions(+), 90 deletions(-) diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 257babb535a..a89b64df276 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -11,7 +11,8 @@ use lance_index::{ optimize::OptimizeOptions, progress::NoopIndexBuildProgress, scalar::{ - CreatedIndex, OldIndexDataFilter, inverted::InvertedIndex, lance_format::LanceIndexStore, + CreatedIndex, OldIndexDataFilter, ScalarIndex, inverted::InvertedIndex, + lance_format::LanceIndexStore, }, }; use lance_select::{RowAddrTreeMap, RowSetOps}; @@ -115,6 +116,40 @@ async fn load_unindexed_training_data( .await } +/// Build a fresh, canonical (non-sharded) scalar index over `fragment_ids`, +/// reusing `reference_index`'s params and training criteria. +async fn rebuild_scalar_segment( + dataset: &Dataset, + reference_index: &Arc, + field_path: &str, + column_name: &str, + uuid: &str, + fragment_ids: Vec, +) -> Result { + let params = reference_index.derive_index_params()?; + let update_criteria = reference_index.update_criteria(); + let training_data = load_training_data( + dataset, + field_path, + &update_criteria.data_criteria, + None, + true, + Some(fragment_ids), + ) + .await?; + super::scalar::build_scalar_index( + dataset, + column_name, + uuid, + ¶ms, + true, + None, + Some(training_data), + Arc::new(NoopIndexBuildProgress), + ) + .await +} + #[allow(clippy::too_many_arguments)] async fn merge_scalar_indices<'a>( dataset: Arc, @@ -132,25 +167,10 @@ async fn merge_scalar_indices<'a>( )); } - let num_to_merge = match index_type { - IndexType::BTree => options - .num_indices_to_merge - .unwrap_or(1) - .min(old_indices.len()), - _ => { - if old_indices.len() > 1 { - return Err(Error::not_supported(format!( - "Cannot optimize multi-segment {:?} scalar index \ - ({} segments): this index type lacks a k-way \ - segment merge primitive. Drop and recreate the \ - index to consolidate.", - index_type, - old_indices.len(), - ))); - } - 1 - } - }; + let num_to_merge = options + .num_indices_to_merge + .unwrap_or(1) + .min(old_indices.len()); // No new data + ≤1 old selected = rewriting one segment to itself. if unindexed.is_empty() && num_to_merge <= 1 { @@ -189,19 +209,29 @@ async fn merge_scalar_indices<'a>( frag_bitmap |= &effective_old_frags; let new_uuid = Uuid::new_v4(); - let created_index = if selected_old_indices.is_empty() { - // Delta: build from new_data only. Old segments stay in manifest. - let params = reference_index.derive_index_params()?; - let unindexed_fragment_ids: Vec = base_unindexed_bitmap.iter().collect(); - super::scalar::build_scalar_index( + // Scalar Index that expos an N:1 segment-merge primitive reachable without + // rescanning the dataset + let has_segment_merge_primitive = matches!(index_type, IndexType::BTree); + + // Merge new data into the existing segment(s) instead of rebuilding from + // scratch, when both hold: + // - `effective_old_frags`: the selected segments' coverage intersected + // with live fragments is non-empty, i.e. there is old data worth keeping. + // - `has_segment_merge_primitive` (Indices supports N:1 segments merge) OR + // `selected_old_indices.len() == 1` (any scalar type can `update` one). + // Otherwise (e.g. ≥2 selected segments of a type without an N:1 merge + // primitive) the index is rebuilt from scratch over `frag_bitmap`. + let can_merge_segments = !effective_old_frags.is_empty() + && (has_segment_merge_primitive || selected_old_indices.len() == 1); + + let created_index = if !can_merge_segments { + rebuild_scalar_segment( dataset.as_ref(), + &reference_index, + field_path, column_name, &new_uuid.to_string(), - ¶ms, - true, - Some(unindexed_fragment_ids), - None, - Arc::new(NoopIndexBuildProgress), + frag_bitmap.iter().collect(), ) .await? } else { @@ -209,53 +239,27 @@ async fn merge_scalar_indices<'a>( let new_data_stream = load_unindexed_training_data(dataset.as_ref(), field_path, &update_criteria, unindexed) .await?; + let new_store = LanceIndexStore::from_dataset_for_new(&dataset, &new_uuid.to_string())?; + let old_data_filter = + build_old_data_filter(dataset.as_ref(), &effective_old_frags, &deleted_old_frags) + .await?; - if effective_old_frags.is_empty() { - // Stale: every selected segment's coverage was retired. - // Rebuild from new_data; `selected` is returned as - // `removed_indices` so the manifest drops them. - let params = reference_index.derive_index_params()?; - super::scalar::build_scalar_index( - dataset.as_ref(), - column_name, - &new_uuid.to_string(), - ¶ms, - true, - None, - Some(new_data_stream), - Arc::new(NoopIndexBuildProgress), - ) - .await? - } else { - let new_store = LanceIndexStore::from_dataset_for_new(&dataset, &new_uuid.to_string())?; - let old_data_filter = - build_old_data_filter(dataset.as_ref(), &effective_old_frags, &deleted_old_frags) - .await?; - - match index_type { - IndexType::BTree => { - crate::index::scalar::btree::open_and_merge_segments( - dataset.as_ref(), - field_path, - selected_old_indices, - new_data_stream, - &new_store, - old_data_filter, - ) + match index_type { + IndexType::BTree => { + crate::index::scalar::btree::open_and_merge_segments( + dataset.as_ref(), + field_path, + selected_old_indices, + new_data_stream, + &new_store, + old_data_filter, + ) + .await? + } + _ => { + reference_index + .update(new_data_stream, &new_store, old_data_filter) .await? - } - _ => { - debug_assert_eq!( - selected_old_indices.len(), - 1, - "num_to_merge clamp invariant violated: index types \ - without a k-way merge primitive must reach \ - update() with exactly one selected segment" - ); - reference_index - .update(new_data_stream, &new_store, old_data_filter) - .await? - } } } }; @@ -1662,32 +1666,37 @@ mod tests { .unwrap(); let dataset = DatasetBuilder::from_uri(test_uri).load().await.unwrap(); - // Bitmap remains a single-segment logical index. The original UUID - // is replaced (because update() produces a new UUID), but we never - // grow to 2 segments. + // append() (= num_indices_to_merge: Some(0)) is now honored uniformly: + // Bitmap, like BTree, must keep the original segment untouched and add + // exactly one delta segment covering only the appended fragment. let committed = dataset.load_indices_by_name("cat_idx").await.unwrap(); assert_eq!( committed.len(), - 1, - "Bitmap optimize append() must silently merge instead of creating a 2nd segment, \ - got {committed:?}" + 2, + "Bitmap optimize append() must add a delta segment, not merge, got {committed:?}" ); - // The merged segment should cover both fragments (old + new). - let frags: std::collections::BTreeSet = committed[0] + assert!( + committed.iter().any(|idx| idx.uuid == original_uuid), + "append() must preserve the pre-existing segment {original_uuid}, got {committed:?}" + ); + let new_segment = committed + .iter() + .find(|idx| idx.uuid != original_uuid) + .expect("append() must add a new delta segment"); + let new_segment_frags: std::collections::BTreeSet = new_segment .fragment_bitmap .as_ref() - .expect("merged Bitmap should carry fragment coverage") + .expect("delta Bitmap should carry fragment coverage") .iter() .collect(); - assert_eq!(frags, [0u32, 1].into_iter().collect()); - // Sanity: the original UUID was replaced by the merge. - assert_ne!( - committed[0].uuid, original_uuid, - "update() should produce a new UUID, but committed segment still has the original" + assert_eq!( + new_segment_frags, + [1u32].into_iter().collect(), + "the delta segment must cover only the appended fragment" ); // Data correctness: a value that lives only in the appended fragment - // is queryable through the index. + // is queryable through the (now multi-segment) index. let rows = dataset .scan() .filter("category = 'd'")