Skip to content

feat(core): process-wide byte budget for MemWAL merges - #265

Merged
beinan merged 1 commit into
lance-format:mainfrom
beinan:feat/merge-memory-budget
Sep 25, 2026
Merged

beinan merged 1 commit into
lance-format:mainfrom
beinan:feat/merge-memory-budget

Conversation

@beinan

@beinan beinan commented Sep 25, 2026

Copy link
Copy Markdown
Collaborator

Problem

merge_max_bytes bounds how much one merge buffers. Nothing bounded how many merges a process ran at once. Every merge entry point — the worker's flush and cleanup sweepers, the count trigger, the manual /merge-wal route, and the master's fan-out — reserved memory independently.

Once the master began scheduling merges for every store (#264) and fanning out to every worker, a worker could be asked to hold a dozen 1 GiB merge reads at once. Production, 2026-09-25: scaling masters 3→6 OOMKilled 17 of 20 workers within 90 seconds (each had served ~370 merges in the hour); scaling back to 3 still produced a second wave of 6 OOMs 75 minutes later. Per-merge caps were all honoured; their sum wasn't.

Fix

MergeMemoryBudget — one per process, measured in bytes, shared by every merge regardless of trigger.

  • Before reading anything, a merge reserves min(merge_max_bytes, budget) in one atomic acquire (Semaphore::acquire_many, 1 permit = 1 MiB). If the budget is exhausted it waits for a release; it is never rejected. A master fanning out to a busy worker sees a slower worker, not a failure.
  • The reservation only grows mid-read for an oversized first generation, which is folded whole (existing rule). By then the holder has at least as much as any waiter could need, so growth can't form a cycle.
  • A request larger than the whole budget is admitted once it is the sole holder — same lone-oversized rule as BlobBudget — so an oversized generation still makes progress.
  • The RAII MergeReservation travels inside PreparedMerge and is dropped right after merge_prepared_batches consumes the batches, before the manifest drain and directory deletes.

Why this can't deadlock: the full initial reservation is taken before any I/O, so no merge ever holds part of what it needs while waiting for the rest.

Server: ROLLOUT_MERGE_MEMORY_BYTES (default 3 GiB; 0 disables) builds the budget in AppState and threads it into rollout/generic/datagen/context store options. With the total bounded, merge_max_bytes becomes an efficiency knob (fewer, larger commits) rather than the OOM safety valve.

Metric: rollout_merge_budget_wait_seconds (job-scale buckets). A rising tail means merges are queueing on memory — intended under load, not an error.

Verification

  • cargo clippy --workspace --all-targets -D warnings clean.
  • merge_budget::tests: reserve/release, exhausted budget makes the next reserve wait and resume on release, oversized request takes the whole budget and proceeds, grow_to waits for the extra and is a downward no-op.
  • generic_store::tests::concurrent_merges_queue_on_the_shared_memory_budget: two real stores, budget sized for one initial reservation; the second prepare_cleanup_merge blocks until the first commit_prepared_merge, then completes; budget reads 0 at the end.
  • Will be load-tested on the staging fleet with 6 masters before production.

Refs #264, #256.

🤖 Generated with Claude Code

merge_max_bytes bounds how much ONE merge buffers. Nothing bounded how many
merges a process ran at once: every entry point -- the worker's flush and
cleanup sweepers, the count trigger, the manual /merge-wal route, and the
master's fan-out -- reserved memory independently. Once the master began
scheduling merges for every store (lance-format#264) and fanning out to every worker,
a worker could be asked to hold a dozen 1 GiB merge reads at once, and 17
of 20 production workers OOMKilled within ninety seconds; a second wave of
6 followed 75 minutes later at the lower master count.

MergeMemoryBudget is the missing bound, measured in bytes, not merges.
Every merge reserves min(merge_max_bytes, budget) from it in one acquire
before reading anything, grows the reservation only for an oversized first
generation (which is folded whole), and carries the RAII reservation inside
PreparedMerge so it is released the moment the commit has consumed the
batches. A merge that cannot fit WAITS for another to release rather than
failing, so a master fanning out to a busy worker sees a slower worker,
not an error. Taking the full initial reservation atomically means a merge
never holds part of what it needs while waiting, which is what rules out
deadlock; an oversized reservation is admitted once it is the sole holder,
mirroring the blob budget's lone-oversized-request rule.

Server: ROLLOUT_MERGE_MEMORY_BYTES (default 3 GiB, 0 disables) builds one
budget in AppState and threads it into every store kind's options. With
the total bounded, merge_max_bytes becomes an efficiency knob (fewer,
larger commits) rather than the OOM safety valve, and can be raised.

Metrics: rollout_merge_budget_wait_seconds (job-scale buckets).

Tests: unit tests for reserve/release/wait/oversized/grow, and a two-store
test where the second merge blocks on the budget until the first commits.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@beinan
beinan merged commit 04a7a24 into lance-format:main Sep 25, 2026
10 checks passed
beinan added a commit that referenced this pull request Sep 29, 2026
… workers merge (#277)

## 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-wal` requests.

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
plugin `provides_exact_answer()`; otherwise
`create_full_table_joined_stream` reads the **whole base table** into a
DataFusion hash join. Our `id` index was a **ZoneMap**, whose plugin
returns `false`, 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`). New `has_key_btree_index()` says whether the
exact-answer index is present; a ZoneMap under the same name does not
count.
- A `merge-wal` task 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. #273 keeps the index
current after each compaction.
- Master detail string and metrics renamed accordingly;
`master_merge_wal_index_built_total` counts pre-merge builds.

A BTree on a string `id` is 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's `MergeInsertJob::explain_plan` which
path a merge would take. With no index it renders a plan containing
`HashJoin`; **with a ZoneMap it still renders `HashJoin`**; with the
BTree it refuses with the scalar-index `NotSupported` (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`: asserts
`has_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.
- Master suite `--include-ignored` 75/75; `storage_reliability` 10/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](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