Skip to content

fix: stop reads and merges from opening the whole shard at once - #278

Merged
beinan merged 1 commit into
lance-format:mainfrom
beinan:fix/read-amplification-and-index-coverage
Sep 30, 2026
Merged

beinan merged 1 commit into
lance-format:mainfrom
beinan:fix/read-amplification-and-index-coverage

Conversation

@beinan

@beinan beinan commented Sep 30, 2026

Copy link
Copy Markdown
Collaborator

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. fix: the id index must be a BTree, and merge-wal builds it before the workers merge #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

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
beinan merged commit c5f658c into lance-format:main Sep 30, 2026
10 checks passed
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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant