fix: stop reads and merges from opening the whole shard at once - #278
Merged
beinan merged 1 commit intoSep 30, 2026
Merged
Conversation
Three memory amplifiers that took workers to their 32 GiB limit, in the order they were observed: 1. A present-but-stale id BTree. Every merge and compaction appends fragments the index does not cover, and `merge_insert` full-scans exactly those (6.5 GB read to merge 2,046 rows). Before fan-out the master now calls `extend_id_btree_index`, which appends an index delta over `unindexed_fragments` via `optimize_indices`, so the probe stays a probe as the table grows. 2. Reads open every flushed generation pending merge. A point lookup over a 16k-generation shard held ~31 GiB. `pending_generations_max` (server `ROLLOUT_WAL_PENDING_MAX_GENERATIONS`, default 4096) refuses such a read with a recognizable error the server maps to 503 `OVERLOADED`; the merge that fixes the shard is the same merge the read was competing with. The master opens stores uncapped: it observes and merges shards *because* they are behind. 3. Merge fan-out. `MERGE_WAL_CONCURRENCY` bounds tasks per master, and every task hits every worker, so six masters at 4 ran 24 merges on one worker at once, each with its own DataFusion pool outside the merge byte budget. `ROLLOUT_MERGE_CONCURRENCY` (default 4) bounds concurrent merge-wal requests per worker; later ones wait for a slot. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
beinan
added a commit
that referenced
this pull request
Sep 30, 2026
…t one (#279) ## Problem The read cap from #278 (`pending_generations_max`) is checked per writer shard. A read opens every pending generation of **every** shard, so a store with 20 shards at ~1,400 generations each (25,588 total) passes the check and one point lookup still opens all 25k datasets — it OOMKilled a 32 GiB worker in production an hour after #278 shipped. ## Fix Keep the per-shard check (fails fast, before other shards are listed) and additionally refuse once the collected snapshots' total exceeds the cap, with the same `PENDING_GENERATIONS_EXCEEDED` prefix the server maps to `503 OVERLOADED`. Docs updated to say "per shard or in total". ## Verification `pending_generation_cap_counts_every_shard`: cap 3, shards with 2+1 generations read fine; 2+2 is refused with "across 2 shards"; merging one shard brings reads back. Existing per-shard test unchanged; core lib suite passes. 🤖 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
#280) ## Problem The WAL merge uses `WhenMatched::UpdateAll`. Lance's upsert takes every matched target row **with all of its columns** to join and rewrite the fragment. Rollout rows carry multi-megabyte inline blobs, so in production: - a merge matching 17k rows read 7.6 GB (459 KB/row); one matching 1.4k rows read 4.0 GB (3 MB/row) - two such merges in flight held a 32 GiB worker at its limit — cgroup `anon` 27–30 GiB, dropping to 8 GiB within two minutes of the merges finishing - 8 workers OOMKilled in two hours on the #278 image even with merges bounded to two per worker ## Fix `merge_prepared_batches` is now two commits: 1. a **delete-only** `merge_insert` on the key column (`WhenMatched::Delete`, `WhenNotMatched::DoNothing`) — probes the id BTree, touches only the key column and deletion vectors; 2. a plain `Dataset::append` of the new rows into fresh fragments. No target blob is ever read. Last-write-wins is unchanged (`read_flushed_generations` already reduces to one newest row per key). A crash between the commits leaves the old rows deleted and the new rows still in the WAL; the retry deletes nothing and appends once — same end state. Debug builds of the two commits overflow the 2 MiB test-thread stack inside DataFusion's optimizer walk (release builds fit). Tests get 8 MiB via `.cargo/config.toml` `[env] RUST_MIN_STACK` and the same variable in the CI workflow env. ## Verification - New `merge_deletes_old_rows_and_appends_instead_of_rewriting_fragments`: after overwriting one of two keys, the original fragment is kept with a deletion vector and the new row lands in a second fragment; the read returns the newer content exactly once. - Full core suite (247 tests), server, master (incl. etcd `--ignored`) and `lance-context` suites pass; clippy clean. 🤖 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
…on is unchanged (#281) ## Problem The scanner's unchanged-version shortcut reuses the previous stats row wholesale. WAL flushes never bump the base-table version, so a store that no worker or master has ever merged keeps its version forever while its shards fill up. In production `mai3_bigclimb_run6p5_77b_t0r1` sat at "3 pending" in the stats table while its 20 shards accumulated 17,254 generations; the merge sweep reads that row, so it never scheduled a merge, and after #278/#279 every read on the store was refused (503) by the pending-generation cap. ## Fix Both `observe_generic` and `observe_one` still skip the base-table row count on an unchanged version, but recount `pending_wal_generations` (one small shard-manifest read per shard, under the existing observe timeout) and write it into the reused row. ## Verification - `skipped_round_recounts_pending_generations`: 1 → 4 flushed generations with the base version unchanged; the skipped row reports 4. - `generic_store_is_observed_with_pending_wal` extended the same way (3 → 5). - Master suite incl. etcd `--include-ignored` passes; clippy clean. 🤖 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
Oct 1, 2026
…#286) ## Problem #282 builds the key BTree before the delete-only `merge_insert` when the base has **no** index. But an index that exists and covers only some fragments is just as bad: `merge_insert` hash-joins every fragment the index does not cover. Production, 00:55 UTC: `bp-prod2rubric-…-ghcp-100` (2.2 GB base, 65 fragments, BTree built when it had 5) — a merge of 64 generations totalling **37 MB** took a worker from 5 to 32 GiB and OOMKilled 19/20 workers that merged it concurrently. Reproduced deterministically on one worker (5 → 20+ GiB → killed). After a single `optimize_indices` the identical merge on the identical worker peaked at **3 GiB** and finished in 39 s. ## Fix In `merge_prepared_batches`: if the base has fragments and no BTree, build it (unchanged); **otherwise call `extend_key_btree_index`** so the index covers everything appended since — the same thing the master does before fan-out (#278), now also on the worker's own merge path (timer, manual route). New counter `rollout_merge_index_extended_on_demand_total`. ## Verification - New `merge_extends_a_stale_id_btree_over_new_fragments`: after four merges at most one fragment (the one just appended) is uncovered. - `extend_id_btree_index_covers_fragments_added_since_build` updated to the new invariant (1 uncovered after two merges, not 2). - Full core suite passes; clippy clean. 🤖 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
Oct 1, 2026
…slot (#289) ## Problem `ROLLOUT_MERGE_AFTER_GENERATIONS` is the pending-driven merge trigger: when a worker's own shard has ≥ N flushed generations, the 30 s flush sweeper merges it. That is exactly what hot stores need (a store written every few seconds regrows hundreds of generations between the master's 600 s sweeps, and every read has to open all of them — 8–13 s per read on a 176-generation store today). But that sweeper merge bypassed the per-worker merge slot (`ROLLOUT_MERGE_CONCURRENCY`, #278). Enabling the count trigger would stack sweeper merges on top of the master's merge requests with no memory bound. ## Fix `flush_pass` takes the worker's merge-slot semaphore and acquires a permit around `merge_if_due`, so sweeper-initiated and master-initiated merges share one bound. No behaviour change while the count trigger is `0` (the default and current production value). ## Verification Server suite passes (85); clippy clean. Production enablement is a config change in the deployment repo (`ROLLOUT_MERGE_AFTER_GENERATIONS=64`), validated on staging first. 🤖 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
Three memory amplifiers took RocketKeep workers to their 32 GiB limit (25 OOM kills in one 20-minute window this morning):
merge_insertfull-scans exactly those. Observed: 6.5 GB read to merge 2,046 rows.MERGE_WAL_CONCURRENCYis 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
StorageBase::extend_key_btree_index/RolloutStore::extend_id_btree_index: append an index delta overunindexed_fragmentsviaoptimize_indices. The master calls it inensure_id_btree_indexwhen 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.pending_generations_max(serverROLLOUT_WAL_PENDING_MAX_GENERATIONS, default 4096,0disables):wal_shard_snapshotsrefuses a read past the cap withPENDING_GENERATIONS_EXCEEDED, which the server maps to503 OVERLOADED. The master opens stores uncapped (Some(0)) since it observes and merges shards because they are behind.ROLLOUT_MERGE_CONCURRENCY(default 4,0disables): a worker-wide semaphore around merge-wal requests (rollout and generic); later ones wait for a slot.rollout_wal_merge_slot_wait_secondsshows the queueing.Verification
extend_id_btree_index_covers_fragments_added_since_build: 0 without BTree, 0 on full coverage, 2 after two merges thenunindexed_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 toAppError::Overloaded.--include-ignored) pass except two pre-existing#[ignore]benches that also fail onmain.🤖 Generated with Claude Code