From 056b8891ff6e3f8a44ac2abf0f920cd843b6d345 Mon Sep 17 00:00:00 2001 From: Beinan Date: Tue, 29 Sep 2026 22:52:29 +0000 Subject: [PATCH] fix: the id index must be a BTree, and merge-wal builds it before the workers merge Lance's merge_insert joins the WAL rows to the base table on the key. It probes a scalar index only when that index answers equality exactly; otherwise it reads the whole base table into a hash join. Our id index was a ZoneMap, whose plugin reports provides_exact_answer() = false, so every merge of every store did the full join. On a 2 TB, 5.8M-row base table that ran for each 64-generation merge; at 2026-09-29 22:12 UTC a merge sweep of 184 targets put ~24 such joins on every worker and 18 of 20 were OOMKilled together. The merge memory budget (#265) never saw it: it covers reading the WAL, not the join. It also explains why pending generations went from 17k to 426k in a day: merge cost scaled with base-table size, not with the rows merged, so the large stores could not be drained. - create_key_zonemap_index -> create_key_btree_index (IndexType::BTree). has_key_btree_index() reports whether the exact-answer index exists; a ZoneMap under the same name does not count. - A merge-wal task on a rollout store without the BTree builds it first (INDEX_BEFORE_MERGE, default on), so the very next merge probes. This is the fix for the stores already in trouble; #273 keeps the index current after each compaction. Test merge_insert_probes_the_id_btree_instead_of_scanning pins the behaviour from Lance's side: explain_plan renders a HashJoin with no index and with a ZoneMap, and refuses (indexed path) with a BTree. Co-Authored-By: Claude Fable 5 --- .../lance-context-core/src/datagen_store.rs | 2 +- .../lance-context-core/src/generic_store.rs | 2 +- .../lance-context-core/src/rollout_store.rs | 95 +++++++++++++-- crates/lance-context-core/src/store_base.rs | 53 ++++++--- .../tests/storage_reliability.rs | 2 +- crates/lance-context-master/src/config.rs | 6 + crates/lance-context-master/src/routes.rs | 1 + crates/lance-context-master/src/scheduler.rs | 108 +++++++++++++++++- crates/lance-context-master/src/state.rs | 1 + crates/lance-context-master/src/task_store.rs | 1 + 10 files changed, 243 insertions(+), 28 deletions(-) diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 69b4970..1b83ea9 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -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 diff --git a/crates/lance-context-core/src/generic_store.rs b/crates/lance-context-core/src/generic_store.rs index 6cf3fd9..8520025 100644 --- a/crates/lance-context-core/src/generic_store.rs +++ b/crates/lance-context-core/src/generic_store.rs @@ -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 diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index a2fd2fc..6857e46 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -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 { + self.base.has_key_btree_index().await } /// Whether the base table has accumulated at least `min_fragments` @@ -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(); @@ -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 { @@ -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. diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index 148108a..7f67572 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -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()) @@ -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 { + 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 diff --git a/crates/lance-context-core/tests/storage_reliability.rs b/crates/lance-context-core/tests/storage_reliability.rs index c052efa..a81e9e9 100644 --- a/crates/lance-context-core/tests/storage_reliability.rs +++ b/crates/lance-context-core/tests/storage_reliability.rs @@ -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); diff --git a/crates/lance-context-master/src/config.rs b/crates/lance-context-master/src/config.rs index bba4dcd..081dd57 100644 --- a/crates/lance-context-master/src/config.rs +++ b/crates/lance-context-master/src/config.rs @@ -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, diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index d3a3ebf..e0e8eb7 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -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, diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 2ffe22e..73fd151 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -336,10 +336,10 @@ async fn index_id_inner(state: &Arc, name: &str) -> Result, name: &str) -> Result<(), String> { + let uri = state.rollout_uri(name); + let mut store = RolloutStore::open_existing_with_options(&uri, state.rollout_store_options()) + .await + .map_err(|e| e.to_string())?; + if store + .has_id_btree_index() + .await + .map_err(|e| e.to_string())? + { + return Ok(()); + } + let started = std::time::Instant::now(); + store + .create_id_btree_index() + .await + .map_err(|e| format!("building id index before merge: {e}"))?; + metrics::counter!("master_merge_wal_index_built_total").increment(1); + tracing::info!( + target = %name, + elapsed_secs = started.elapsed().as_secs(), + "built id BTree index before merge-wal so merge_insert probes instead of scanning" + ); + Ok(()) +} + async fn run_merge_wal(state: &Arc, target: &str) -> Result { let endpoints = &state.config.worker_endpoints; if endpoints.is_empty() { return Err("no worker endpoints configured (--worker-endpoints)".to_string()); } let (kind, name) = parse_target(target); + if kind == StoreKind::Rollout && state.config.index_before_merge { + ensure_id_btree_index(state, name).await?; + } let route = match kind { StoreKind::Rollout => "api/v1/internal/merge-wal", StoreKind::Generic => "api/v1/generic", @@ -871,6 +910,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: 2, @@ -1009,7 +1049,7 @@ mod tests { assert_eq!(index.depends_on, vec![compact.id.clone()]); let status = await_terminal(&state, &index.id).await; assert_eq!(status.state, TaskState::Done, "got {status:?}"); - assert_eq!(status.detail.as_deref(), Some("built zonemap index on id")); + 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(); @@ -1048,6 +1088,66 @@ mod tests { assert!(!is_missing_fragment_error("HTTP 500 Internal Server Error")); } + /// A merge-wal on a store with no id BTree builds one before fanning out, + /// so the workers' `merge_insert` probes instead of scanning; a second + /// merge finds it present and builds nothing. + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn merge_wal_builds_the_id_btree_first() { + use axum::{routing::post, Json, Router}; + let app = Router::new().route( + "/api/v1/internal/merge-wal/{name}", + post(|| async { Json(serde_json::json!({ "reclaimed": 0 })) }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + let dir = TempDir::new().unwrap(); + let mut cfg = config(&dir); + cfg.worker_endpoints = vec![format!("http://{addr}")]; + cfg.index_before_merge = true; + let state = MasterState::new(cfg).await.unwrap(); + let worker = spawn_scheduler(&state); + + let name = "exp"; + let uri = state.rollout_uri(name); + { + let mut store = RolloutStore::open(&uri).await.unwrap(); + store.add(&[rollout_record("r0")]).await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + assert!(!store.has_id_btree_index().await.unwrap()); + } + state + .registry + .write() + .await + .upsert(name, &uri) + .await + .unwrap(); + + let rec = enqueue(&state, TaskKind::MergeWal, name).await.unwrap(); + assert_eq!(await_terminal(&state, &rec.id).await.state, TaskState::Done); + let store = RolloutStore::open_existing_with_options(&uri, Default::default()) + .await + .unwrap(); + assert!( + store.has_id_btree_index().await.unwrap(), + "merge built the index" + ); + let version_after_first = store.version(); + + // Present now: the next merge does not rebuild it. + let rec = enqueue(&state, TaskKind::MergeWal, name).await.unwrap(); + assert_eq!(await_terminal(&state, &rec.id).await.state, TaskState::Done); + let store = RolloutStore::open_existing_with_options(&uri, Default::default()) + .await + .unwrap(); + assert_eq!(store.version(), version_after_first); + + worker.abort(); + } + /// A base table whose manifest names a missing data file fails every /// compaction with `Not found`. That failure enqueues a `Repair` and the /// compaction again behind it; the repair drops the dead fragment and is @@ -1164,7 +1264,7 @@ mod tests { let rec = enqueue(&state, TaskKind::IndexId, name).await.unwrap(); let status = await_terminal(&state, &rec.id).await; assert_eq!(status.state, TaskState::Done, "got {status:?}"); - assert_eq!(status.detail.as_deref(), Some("built zonemap index on id")); + assert_eq!(status.detail.as_deref(), Some("built btree index on id")); worker.abort(); } diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index b3dadbf..ede5157 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -253,6 +253,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, diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index d8296e3..63ef896 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -1314,6 +1314,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,