diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index fc9f3cb..e343189 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -4103,6 +4103,45 @@ mod tests { /// Observable shape: an overwritten key leaves the old fragment in /// place with a deletion vector, and the new row lands in a fresh /// fragment; last-write-wins still holds. + /// A BTree that exists but covers only some fragments leaves the rest to + /// a full scan in the delete-only merge_insert. Each merge extends it + /// over whatever the previous merges and compactions appended, so no + /// fragment is ever scanned twice. + #[test] + fn merge_extends_a_stale_id_btree_over_new_fragments() { + use lance::index::DatasetIndexInternalExt as _; + let dir = TempDir::new().unwrap(); + let uri = dir.path().to_string_lossy().to_string(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let mut store = RolloutStore::open(&uri).await.unwrap(); + store.add(&[assistant_record("a-0")]).await.unwrap(); + store.flush().await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + // Three more merges, each appending a fragment the index (built by + // the first merge) does not cover. + for id in ["a-1", "a-2", "a-3"] { + store.add(&[assistant_record(id)]).await.unwrap(); + store.flush().await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + } + let unindexed = store + .base + .dataset + .unindexed_fragments(ROLLOUT_ID_INDEX_NAME) + .await + .unwrap() + .len(); + // The merge that appended the last fragment extended the index + // before its own write, so at most that one fragment is uncovered. + assert!( + unindexed <= 1, + "unindexed fragments after merges: {unindexed}" + ); + assert_eq!(store.list(None, None).await.unwrap().len(), 4); + }); + } + #[test] fn merge_deletes_old_rows_and_appends_instead_of_rewriting_fragments() { let dir = TempDir::new().unwrap(); @@ -4398,7 +4437,9 @@ mod tests { store.create_id_btree_index().await.unwrap(); assert_eq!(store.extend_id_btree_index().await.unwrap(), 0); - // Two merges land two fragments the index knows nothing about. + // Two merges land two fragments. Each merge extends the index over + // what was uncovered *before* its own append, so only the last + // appended fragment is left for the explicit call here. for id in ["a-1", "a-2"] { store.add(&[assistant_record(id)]).await.unwrap(); store.flush().await.unwrap(); @@ -4414,9 +4455,9 @@ mod tests { .len() } }; - assert_eq!(unindexed(&store).await, 2); + assert_eq!(unindexed(&store).await, 1); - assert_eq!(store.extend_id_btree_index().await.unwrap(), 2); + assert_eq!(store.extend_id_btree_index().await.unwrap(), 1); assert_eq!(unindexed(&store).await, 0); assert!(store.has_id_btree_index().await.unwrap()); assert_eq!(store.extend_id_btree_index().await.unwrap(), 0); diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index bc1b1b9..f1b5f2d 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -1232,13 +1232,26 @@ impl StorageBase { // that tried to merge it; building the BTree first took 19 s. The // master builds it before fan-out (#277), but a worker's own timer // or a manual merge must not depend on the master having been here. - if self.dataset.count_fragments() > 0 && !self.has_key_btree_index().await? { - metrics::counter!("rollout_merge_index_built_on_demand_total").increment(1); - info!( - uri = %self.dataset.uri(), - "base table has no key BTree; building it before the merge" - ); - self.create_key_btree_index().await?; + if self.dataset.count_fragments() > 0 { + if !self.has_key_btree_index().await? { + metrics::counter!("rollout_merge_index_built_on_demand_total").increment(1); + info!( + uri = %self.dataset.uri(), + "base table has no key BTree; building it before the merge" + ); + self.create_key_btree_index().await?; + } else { + // Present is not enough: every merge and compaction appends + // fragments the index does not cover, and the delete-only + // merge_insert hash-joins exactly those. A 2.2 GB store whose + // BTree covered 1 of 65 fragments took a worker from 5 to + // 32 GiB on a 37 MB merge and OOMKilled 19 of them at once; + // after one `optimize_indices` the same merge peaked at 3 GiB. + let covered = self.extend_key_btree_index().await?; + if covered > 0 { + metrics::counter!("rollout_merge_index_extended_on_demand_total").increment(1); + } + } } let key_index = merge_schema.index_of(&self.key_column)?; let key_schema = Arc::new(merge_schema.project(&[key_index])?); @@ -1574,9 +1587,14 @@ impl StorageBase { if unindexed == 0 { return Ok(0); } + // `merge(1)`, not `append()`: fold the new fragments into the existing + // delta so the index stays a single delta. Compaction refuses to bin + // fragments covered by different index-delta sets together, so one + // delta per merge would leave every fragment in its own bin and + // compaction would never coalesce anything. self.dataset .optimize_indices( - &OptimizeOptions::append().index_names(vec![ID_INDEX_NAME.to_string()]), + &OptimizeOptions::merge(1).index_names(vec![ID_INDEX_NAME.to_string()]), ) .await?; self.reload().await?;