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
2 changes: 1 addition & 1 deletion crates/lance-context-core/src/datagen_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -423,7 +423,7 @@ impl DatagenStore {
/// Idempotent. Datagen previously had no scalar index, so every point
/// lookup by event id scanned.
pub async fn create_event_id_index(&mut self) -> LanceResult<()> {
self.base.create_key_zonemap_index().await
self.base.create_key_btree_index().await
}

/// Seal the active memtable. A no-op here in normal operation, since
Expand Down
2 changes: 1 addition & 1 deletion crates/lance-context-core/src/generic_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -453,7 +453,7 @@ impl GenericStore {

/// Build a ZoneMap scalar index on `id`. Idempotent.
pub async fn create_id_index(&mut self) -> LanceResult<()> {
self.base.create_key_zonemap_index().await
self.base.create_key_btree_index().await
}

/// Row count of the base table. Excludes rows still in unmerged
Expand Down
95 changes: 88 additions & 7 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -702,9 +702,14 @@ impl RolloutStore {
}

/// Build a ZoneMap scalar index on the base table's `id` column. Idempotent.
/// See `StorageBase::create_key_zonemap_index`.
pub async fn create_id_zonemap_index(&mut self) -> LanceResult<()> {
self.base.create_key_zonemap_index().await
/// See `StorageBase::create_key_btree_index`.
pub async fn create_id_btree_index(&mut self) -> LanceResult<()> {
self.base.create_key_btree_index().await
}

/// See `StorageBase::has_key_btree_index`.
pub async fn has_id_btree_index(&self) -> LanceResult<bool> {
self.base.has_key_btree_index().await
}

/// Whether the base table has accumulated at least `min_fragments`
Expand Down Expand Up @@ -4196,9 +4201,80 @@ mod tests {
});
}

/// The whole point of the id index: `merge_insert` must take the indexed
/// probe path, not a full-table hash join. Lance's `explain_plan` only
/// renders the full-scan plan and returns `NotSupported` when the job
/// would use a scalar index, so "explain refuses" is the observable
/// signal that the merge will probe. Without the index (and with a
/// ZoneMap, which cannot answer equality exactly) explain succeeds and
/// shows a HashJoin over a LanceScan of the base table.
#[test]
fn merge_insert_probes_the_id_btree_instead_of_scanning() {
use lance::dataset::MergeInsertBuilder;
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();

let plan_for = |s: &RolloutStore| {
let dataset = Arc::new(s.base.dataset.clone());
async move {
let mut b =
MergeInsertBuilder::try_new(dataset, vec!["id".to_string()]).unwrap();
b.when_matched(lance::dataset::WhenMatched::UpdateAll);
b.try_build().unwrap().explain_plan(None, false).await
}
};

// No index: full-table join.
let plan = plan_for(&store).await.expect("full-scan plan renders");
assert!(plan.contains("HashJoin"), "{plan}");

// ZoneMap: still a full-table join (not an exact-answer index).
store
.base
.dataset
.create_index_builder(
&["id"],
lance_index::IndexType::ZoneMap,
&lance_index::scalar::ScalarIndexParams::default(),
)
.name(ROLLOUT_ID_INDEX_NAME.to_string())
.replace(true)
.await
.unwrap();
store.base.reload().await.unwrap();
assert!(!store.has_id_btree_index().await.unwrap());
let plan = plan_for(&store)
.await
.expect("ZoneMap does not change the plan");
assert!(plan.contains("HashJoin"), "{plan}");

// BTree: indexed path, which explain_plan cannot render.
store.create_id_btree_index().await.unwrap();
assert!(store.has_id_btree_index().await.unwrap());
let err = plan_for(&store).await.expect_err("indexed path");
assert!(
err.to_string().contains("scalar-index"),
"expected the scalar-index refusal, got: {err}"
);

// And a real merge through the index still works.
store.add(&[assistant_record("a-0")]).await.unwrap();
store.add(&[assistant_record("a-1")]).await.unwrap();
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();
assert_eq!(store.list(None, None).await.unwrap().len(), 2);
});
}

#[test]
fn create_id_zonemap_index_builds_and_is_idempotent() {
// Building the ZoneMap index on `id` must succeed even though the
fn create_id_btree_index_builds_and_is_idempotent() {
// Building the BTree index on `id` must succeed even though the
// rollout table also carries a (fieldless) MemWAL index, and calling it
// twice must not error (replace(true) rebuilds in place).
let dir = TempDir::new().unwrap();
Expand All @@ -4213,7 +4289,12 @@ mod tests {
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();

store.create_id_zonemap_index().await.unwrap();
assert!(!store.has_id_btree_index().await.unwrap());
store.create_id_btree_index().await.unwrap();
assert!(
store.has_id_btree_index().await.unwrap(),
"the id index must be a BTree: merge_insert only probes an exact-answer index"
);
let has_id_index = |s: &RolloutStore| {
let dataset = s.base.dataset.clone();
async move {
Expand All @@ -4228,7 +4309,7 @@ mod tests {
assert!(has_id_index(&store).await, "id index should exist");

// Idempotent: a second build replaces in place without erroring.
store.create_id_zonemap_index().await.unwrap();
store.create_id_btree_index().await.unwrap();
assert!(has_id_index(&store).await, "id index should still exist");

// Rows remain readable exactly once after indexing.
Expand Down
53 changes: 39 additions & 14 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1434,28 +1434,39 @@ impl StorageBase {
})
}

/// Build a ZoneMap scalar index on the base table's key column.
/// Build a BTree scalar index on the base table's key column, or report
/// whether one already exists.
///
/// The key column is the table's (unenforced) primary key, so a lightweight
/// per-fragment min/max index accelerates point lookups and range scans on
/// the already-flushed base table. `replace(true)` makes this idempotent.
/// The key column is the table's (unenforced) primary key. The index does
/// two jobs:
///
/// - point lookups and range scans on the already-merged base table;
/// - **the MemWAL merge itself**. `merge_insert` joins the WAL rows to the
/// base table on the key. With a scalar index that answers equality
/// exactly it probes only the fragments holding the source keys; without
/// one it reads the *whole* base table into a hash join. On a 2 TB /
/// 5.8M-row table that full join ran for every 64-generation merge and
/// took every worker with it (18/20 OOMKilled together, 2026-09-29).
///
/// It must be a BTree: Lance's `merge_insert` only takes the indexed path
/// for an index whose plugin `provides_exact_answer()`, and ZoneMap (a
/// per-fragment min/max) does not, so a ZoneMap here changed nothing for
/// the merge. `replace(true)` makes rebuilding idempotent.
///
/// # MemWAL interaction
///
/// The base table carries a fieldless MemWAL index, and Lance's MemWAL does
/// not *maintain* ZoneMap indices across WAL flushes (it only keeps the
/// indices named in `maintained_indexes`). That does not affect correctness:
/// rows are de-duplicated by the key column at read time, so the ZoneMap
/// only ever needs to describe the base table's already-merged fragments —
/// rows still living in unmerged WAL generations are found by the normal
/// scan of those generations.
pub async fn create_key_zonemap_index(&mut self) -> LanceResult<()> {
/// Lance's MemWAL does not maintain this index across WAL flushes (it only
/// keeps the indices named in `maintained_indexes`). That does not affect
/// correctness: rows are de-duplicated by the key column at read time, so
/// the index only ever needs to describe the base table's already-merged
/// fragments, and compaction rebuilds it (see the master's `IndexId`).
pub async fn create_key_btree_index(&mut self) -> LanceResult<()> {
self.ensure_writable()?;
info!(column = %self.key_column, "creating ZoneMap index on key column");
info!(column = %self.key_column, "creating BTree index on key column");
self.dataset
.create_index_builder(
&[self.key_column.as_str()],
IndexType::ZoneMap,
IndexType::BTree,
&ScalarIndexParams::default(),
)
.name(ID_INDEX_NAME.to_string())
Expand All @@ -1466,6 +1477,20 @@ impl StorageBase {
self.reload().await
}

/// Whether the base table's key column has a BTree index (the one
/// `merge_insert` can use). A ZoneMap under the same name does not count.
pub async fn has_key_btree_index(&self) -> LanceResult<bool> {
let indices = self.dataset.load_indices().await?;
for index in indices.iter().filter(|i| i.name == ID_INDEX_NAME) {
if let Some(details) = &index.index_details {
if details.type_url.ends_with("BTreeIndexDetails") {
return Ok(true);
}
}
}
Ok(false)
}

/// Whether the base table has accumulated at least `min_fragments`
/// fragments (and is thus worth compacting). Quiet-hours gating from
/// [`CompactionConfig`] is honored so an external scheduler can pass the
Expand Down
2 changes: 1 addition & 1 deletion crates/lance-context-core/tests/storage_reliability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,7 @@ async fn pinned_rollout_rejects_writes_until_refresh() {
.to_string()
.contains("read-only"));
assert!(store.compact(None).await.is_err());
assert!(store.create_id_zonemap_index().await.is_err());
assert!(store.create_id_btree_index().await.is_err());
assert!(store.is_version_pinned());
store.refresh_latest().await.unwrap();
assert_eq!(store.version(), version);
Expand Down
6 changes: 6 additions & 0 deletions crates/lance-context-master/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,12 @@ pub struct MasterConfig {
#[arg(long, env = "INDEX_AFTER_COMPACTION", default_value_t = true, action = clap::ArgAction::Set)]
pub index_after_compaction: bool,

/// Before fanning a merge-wal out to the workers, build the `id` BTree
/// index on the target's base table if it is missing. Without it Lance's
/// `merge_insert` full-scans the base table on every merge.
#[arg(long, env = "INDEX_BEFORE_MERGE", default_value_t = true, action = clap::ArgAction::Set)]
pub index_before_merge: bool,

/// Maximum bytes per compacted output file. `0` uses Lance's default.
#[arg(
long,
Expand Down
1 change: 1 addition & 0 deletions crates/lance-context-master/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -695,6 +695,7 @@ mod tests {
compaction_batch_size: 8,
compaction_max_source_fragments: 32,
index_after_compaction: false,
index_before_merge: false,
compaction_max_bytes_per_file: 1024 * 1024 * 1024,
merge_wal_interval_secs: 0,
merge_wal_min_generations: 8,
Expand Down
Loading
Loading