Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 31 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
9 changes: 9 additions & 0 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1226,6 +1226,15 @@ impl StorageBase {
batches: Vec<RecordBatch>,
merge_schema: Arc<Schema>,
) -> 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
Expand Down
24 changes: 17 additions & 7 deletions crates/lance-context-master/src/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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();
}
Expand Down
Loading