diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index e8095a1..fc9f3cb 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -4066,6 +4066,37 @@ mod tests { }); } + /// A merge into a base table that has rows but no id BTree builds the + /// BTree first: without it the delete-only `merge_insert` is a hash + /// join over the whole base table. The very first merge (empty base) + /// has nothing to join and builds nothing. + #[test] + fn merge_builds_the_id_btree_when_the_base_has_rows_but_no_index() { + 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(); + assert_eq!(store.base.dataset.count_fragments(), 1); + assert!( + !store.has_id_btree_index().await.unwrap(), + "empty base: nothing to index" + ); + + store.add(&[assistant_record("a-1")]).await.unwrap(); + store.flush().await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + assert!( + store.has_id_btree_index().await.unwrap(), + "second merge joins against a populated base and builds the BTree first" + ); + assert_eq!(store.list(None, None).await.unwrap().len(), 2); + }); + } + /// The WAL merge is a delete-only `merge_insert` on the key followed by /// an append: an upsert would take the matched target rows with every /// column (multi-megabyte inline blobs included) to rewrite them. diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index 6f1e32c..b3f10f2 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -1226,6 +1226,15 @@ impl StorageBase { batches: Vec, merge_schema: Arc, ) -> LanceResult<()> { + // Without an exact-answer key index the delete-only merge_insert is + // a hash join over the whole base table. On a 174 GB / 67k-row store + // that no master had ever indexed, that join OOMKilled every worker + // 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? { + self.create_key_btree_index().await?; + } let key_index = merge_schema.index_of(&self.key_column)?; let key_schema = Arc::new(merge_schema.project(&[key_index])?); let keys = batches diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index c9c028f..0587206 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -1068,12 +1068,22 @@ mod tests { assert_eq!(status.state, TaskState::Done, "got {status:?}"); assert_eq!(status.detail.as_deref(), Some("built btree index on id")); - // A second compaction with nothing to rewrite does not enqueue another. - let again = enqueue(&state, TaskKind::Compact, name).await.unwrap(); - assert_eq!( - await_terminal(&state, &again.id).await.state, - TaskState::Done - ); + // Compact until nothing is left to rewrite. The first merge built the + // id BTree over the first fragment, and Lance compacts indexed and + // unindexed fragments in separate groups, so reaching one fragment + // can take more than one pass. Only passes that rewrote fragments + // enqueue an IndexId; the final no-op pass must not. + let mut rewriting_compactions = 1; + loop { + let again = enqueue(&state, TaskKind::Compact, name).await.unwrap(); + let status = await_terminal(&state, &again.id).await; + assert_eq!(status.state, TaskState::Done, "got {status:?}"); + if status.detail.as_deref() == Some("removed 0 / added 0 fragments") { + break; + } + rewriting_compactions += 1; + assert!(rewriting_compactions <= 4, "compaction never converged"); + } let indexes = state .task_store .list() @@ -1082,7 +1092,7 @@ mod tests { .into_iter() .filter(|task| task.kind == TaskKind::IndexId) .count(); - assert_eq!(indexes, 1); + assert_eq!(indexes, rewriting_compactions); worker.abort(); }