fix: the id index must be a BTree, and merge-wal builds it before the workers merge - #277
Merged
Merged
Conversation
… 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 (lance-format#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; lance-format#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 <noreply@anthropic.com>
beinan
added a commit
that referenced
this pull request
Sep 30, 2026
## Problem Three memory amplifiers took RocketKeep workers to their 32 GiB limit (25 OOM kills in one 20-minute window this morning): 1. **Stale id BTree.** #277 builds the BTree before merge, but every merge and compaction afterwards appends fragments it does not cover, and `merge_insert` full-scans exactly those. Observed: 6.5 GB read to merge 2,046 rows. 2. **Reads open every pending generation.** A single point lookup over a shard with 16k flushed generations held ~31 GiB. The only cure is the merge that is already behind, which the read is competing with. 3. **Merge fan-out.** `MERGE_WAL_CONCURRENCY` is per master and every task hits every worker, so six masters at 4 put 24 concurrent merges on each worker, each with its own DataFusion pool that the merge byte budget does not count. ## Fix 1. `StorageBase::extend_key_btree_index` / `RolloutStore::extend_id_btree_index`: append an index delta over `unindexed_fragments` via `optimize_indices`. The master calls it in `ensure_id_btree_index` when the BTree already exists, so the probe stays a probe as the table grows. No-op on full coverage; leaves tables without the BTree alone. 2. `pending_generations_max` (server `ROLLOUT_WAL_PENDING_MAX_GENERATIONS`, default 4096, `0` disables): `wal_shard_snapshots` refuses a read past the cap with `PENDING_GENERATIONS_EXCEEDED`, which the server maps to `503 OVERLOADED`. The master opens stores uncapped (`Some(0)`) since it observes and merges shards *because* they are behind. 3. `ROLLOUT_MERGE_CONCURRENCY` (default 4, `0` disables): a worker-wide semaphore around merge-wal requests (rollout and generic); later ones wait for a slot. `rollout_wal_merge_slot_wait_seconds` shows the queueing. ## Verification - `extend_id_btree_index_covers_fragments_added_since_build`: 0 without BTree, 0 on full coverage, 2 after two merges then `unindexed_fragments == 0`, rows still merge correctly. - `reads_are_refused_past_the_pending_generation_cap`: passes at the cap, refused one past it, resumes after the merge drains. - `pending_generation_cap_is_overloaded`: error maps to `AppError::Overloaded`. - Full workspace suites with local etcd (`--include-ignored`) pass except two pre-existing `#[ignore]` benches that also fail on `main`. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
beinan
added a commit
that referenced
this pull request
Sep 30, 2026
…table (#282) ## Problem The delete-only `merge_insert` from #280 probes the key BTree. When the base table has none, Lance falls back to a hash join over the entire base table. `mai3_bigclimb_run6p5_77b_t0r1` (174 GB, 67k rows, 2.6 MB/row) had never been indexed because the master's stats row for it was stale (#281), so every worker that ran a merge on it — the master's fan-out or a manual `/internal/merge-wal` call — was OOMKilled within ~15 s. Building the BTree first took 19 s and the same merge then peaked at 9 GiB. ## Fix `StorageBase::merge_prepared_batches` builds the key BTree (`create_key_btree_index`, idempotent) when the base table has fragments and `has_key_btree_index()` is false, before the delete-only merge. An empty base table has nothing to join and is left unindexed. The master's pre-fan-out build (#277) stays; this is the worker-side guarantee that does not depend on the master having visited the store. ## Verification - `merge_builds_the_id_btree_when_the_base_has_rows_but_no_index`: first merge into an empty base builds nothing; the second merge, against a populated base, creates the BTree and rows still read back exactly once. - Full core suite (249 tests) passes; clippy clean across the workspace. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
At 2026-09-29 22:12 UTC, 18 of 20 production workers were OOMKilled within 66 seconds. The merge sweep had queued 184 targets (normal: 20–50) and every worker was running ~24 concurrent
merge-walrequests.The memory was not the WAL rows (the #265 budget bounds those). It was
merge_insert. Lance's merge_insert joins the source rows to the base table on the key, and it probes a scalar index only when that index's pluginprovides_exact_answer(); otherwisecreate_full_table_joined_streamreads the whole base table into a DataFusion hash join. Ouridindex was a ZoneMap, whose plugin returnsfalse, so every merge of every store did the full join. The largest affected store has a 1.95 TB, 5.8M-row base table: each 64-generation merge (a few hundred rows) read 2 TB.That also explains the trend line: pending generations went from 17k to 426k in 24 h, 40 stores over 4k pending. Merge cost scaled with base-table size rather than with the rows merged, so the large stores could not be drained, so more merges queued, so more full joins ran concurrently.
#273 (index after compaction) did not help because it built a ZoneMap.
Fix
create_key_zonemap_index→create_key_btree_index(IndexType::BTree). Newhas_key_btree_index()says whether the exact-answer index is present; a ZoneMap under the same name does not count.merge-waltask on a rollout store without the BTree builds it before fanning out to the workers (INDEX_BEFORE_MERGE, default on; one manifest read when already present). This is what rescues the stores already in trouble: the very next merge probes. feat(master): rebuild the id ZoneMap index after every compaction that rewrote fragments #273 keeps the index current after each compaction.master_merge_wal_index_built_totalcounts pre-merge builds.A BTree on a string
idis larger than a ZoneMap (it stores every key), tens of MB for 5.8M rows. That is the cost of a merge that reads MBs instead of TBs.Verification
merge_insert_probes_the_id_btree_instead_of_scanning(core): builds a store, then asks Lance'sMergeInsertJob::explain_planwhich path a merge would take. With no index it renders a plan containingHashJoin; with a ZoneMap it still rendersHashJoin; with the BTree it refuses with the scalar-indexNotSupported(explain only renders the full-scan plan), which is the observable signal for the indexed path. A real merge through the index then succeeds.create_id_btree_index_builds_and_is_idempotent: assertshas_id_btree_index()false before / true after, idempotent rebuild.merge_wal_builds_the_id_btree_first(master, etcd): a merge on an unindexed store leaves a BTree behind; a second merge does not bump the version.--include-ignored75/75;storage_reliability10/10; clippy-D warnings; fmt.Rollout note
Prod is holding on
MERGE_WAL_CONCURRENCY=1(env) since the incident. Once this ships, large stores get a BTree on their first merge, after which merge memory is independent of base-table size and the concurrency can go back up.🤖 Generated with Claude Code