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
47 changes: 44 additions & 3 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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();
Expand All @@ -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);
Expand Down
34 changes: 26 additions & 8 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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])?);
Expand Down Expand Up @@ -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?;
Expand Down
Loading