Skip to content

fix(server): stop generic-store WAL merges starving behind rollout sweeps - #256

Merged
beinan merged 2 commits into
lance-format:mainfrom
beinan:fix/generic-wal-merge-starvation
Sep 21, 2026
Merged

beinan merged 2 commits into
lance-format:mainfrom
beinan:fix/generic-wal-merge-starvation

Conversation

@beinan

@beinan beinan commented Sep 21, 2026

Copy link
Copy Markdown
Collaborator

Problem

The global cleanup sweeper (spawn_global_sweeper) walks store kinds serially: every resident rollout store first (each with a 5 × ROLLOUT_CLEANUP_INTERVAL_SECS timeout), then datagen, then generic. On a production worker with hundreds of resident rollout stores and frequent Too many concurrent writers / timeouts, the generic pass was never reached — 48h of logs contain zero kind="generic" sweeper lines.

Meanwhile the 30s flush sweeper runs the count-triggered merge for rollout stores only; for generic stores it just seals. So a hot generic store had no working merge path at all.

Result: one generic store (~20 shards, one-row-per-append traffic) accumulated ~12.6k unmerged flushed generations. GenericStore::get/list → lsm_scanner() → wal_shard_snapshots() unions every flushed generation, so each read re-opened all 12.6k as separate datasets and workers OOMKilled at 32 GiB within ~6 minutes of restart (20 workers, 20–32 restarts each over 3 days).

Manually invoking POST /generic/{name}/merge-wal also showed the generic merge holds the store's write lock for the entire generation read, blocking all reads/appends to that store for the duration (flush sweeper timed out ... kind="generic").

Fix

Bring generic stores to parity with rollout on the maintenance path:

  • state.rs: run the rollout/datagen/generic cleanup passes with tokio::join! so no kind is queued behind another.
  • sweeper.rs: the generic flush() also calls maybe_merge_wal() (count trigger, ROLLOUT_MERGE_AFTER_GENERATIONS), matching rollout.
  • sweeper.rs + generic_store.rs: the generic cleanup merge uses the existing StorageBase prepare/commit split — read lock for the expensive generation read, write lock only for the commit. Adds GenericStore::prepare_cleanup_merge / commit_prepared_merge as thin passthroughs (same as RolloutStore).

No format or default-config changes.

Verification

  • Log analysis on the crashed workers (99% of final-minute lines are _mem_wal/.../gen_NNN loads of the single generic store; 20 shards × 900–1550 pending generations).
  • Mitigation in progress on the live fleet via the existing merge-wal endpoint; the fix is what stops recurrence.
  • CI covers the existing sweeper/generic tests.

🤖 Generated with Claude Code

Beinan Wang and others added 2 commits September 21, 2026 17:30
…eeps

The global cleanup sweeper walked store kinds serially: every resident
rollout store first, each guarded by a multi-minute timeout, then datagen,
then generic. On a worker with hundreds of resident rollout stores the
generic pass was effectively never reached, so a hot generic store's
flushed generations were never folded into its base table. One production
store accumulated ~12.6k unmerged generations across 20 shards; because
the LSM read path unions every flushed generation, each read re-opened all
of them and the worker exceeded its 32 GiB limit within minutes.

Three changes:
- The cleanup sweeper runs the rollout/datagen/generic merge passes
  concurrently, so no kind waits behind another.
- The generic flush sweeper also runs the count-triggered merge
  (maybe_merge_wal), matching rollout, so a hot generic store is bounded
  by ROLLOUT_MERGE_AFTER_GENERATIONS even between cleanup ticks.
- The generic cleanup merge uses the prepare/commit split (read lock for
  the expensive generation read, write lock only for the commit) instead
  of holding the exclusive lock for the whole merge, so a long merge no
  longer blocks the store's appends and reads.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@beinan
beinan merged commit 897a536 into lance-format:main Sep 21, 2026
10 checks passed
beinan added a commit that referenced this pull request Sep 22, 2026
…261)

Closes #259.

## Problem

`wal_shard_snapshots()` hands every flushed generation of every shard to
the `LsmScanner`, and each one is opened as its own dataset on every
`get`/`list`. There is no cap, no warning and no metric on that count.
One generic store reached ~12.6k pending generations across 20 shards
before anyone noticed — by then every read exhausted a 32 GiB worker and
the fleet OOM-looped for three days. The store was only identified by
reading crash logs by hand (99% of the final-minute lines were
`_mem_wal/<shard>/<gen>` loads).

## Fix

On every LSM read, per shard:
- record `flushed_generations.len()` into a new
`rollout_wal_pending_generations` histogram (unlabelled, per the
cardinality rules in `metrics.rs`);
- `warn!` with `uri`, `shard`, `pending`, `warn_at` when the count
reaches `ROLLOUT_WAL_PENDING_WARN_GENERATIONS` (default 256; `0`
disables the warn, the metric is always emitted).

256 sits well above the healthy band (`MERGE_AFTER_GENERATIONS` 50,
`MERGE_MAX_GENERATIONS` 8) and far below the counts that make a read
unaffordable (~1,500/shard).

Plumbing: `pending_generations_warn: Option<usize>` on
`StorageBaseOptions` and every public store options struct
(rollout/generic/datagen/context), server config flag + env,
`lance-context-metrics` gives the histogram count-scale buckets (8, 16,
…, 4096) since the `_duration_seconds` suffix rule does not apply.

No behaviour change.

## Verification

- `cargo clippy --workspace --all-targets -D warnings` clean with and
without the `metrics` feature.
- New test
`generic_store::tests::reads_sample_pending_generations_per_shard`
asserts a `list()` after three sealed adds samples `3` on the histogram
via `DebuggingRecorder`.
- Will be deployed to the staging worker set and confirmed against the
store from #256 before the follow-ups (#257, #258) land.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Beinan Wang <>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
beinan added a commit that referenced this pull request Sep 22, 2026
…it lock split (#263)

Closes #257.

## Problem

`Sweepable::flush` for rollout and generic sealed under the read lock,
then took the **write** lock and called `maybe_merge_own_shard()` /
`maybe_merge_wal()`. Those read up to `MERGE_MAX_GENERATIONS`
generations from object storage *before* committing, so the exclusive
lock was held across the whole read — seconds to tens of seconds on
ABFS. Every `add`/`get`/`list` on the store stalled, and the next 30s
flush tick queued behind it.

Production after #256: `flush sweeper timed out` 52× in 15 min on one
hot generic store, request handlers failing with `Too many concurrent
writers … Attempted 5 times`. Warnings went to zero once the backlog
drained, confirming hold time tracks the merge read.

A second, smaller bug in the same function: a merge failure after a
successful seal was returned as the flush's error, so it was logged as
`flush sweeper failed` and counted under
`rollout_wal_flush_total{result="failed"}`.

## Fix

Apply the prepare (read lock) / commit (write lock) split — already used
by the cleanup path — to the count trigger:

- `StorageBase::prepare_count_merge()`: shared-lock half of
`maybe_merge_own_shard`; `None` when the trigger is disabled
(`merge_after_generations == 0`) or not met. Thin passthroughs on
`RolloutStore` and `GenericStore`, mirroring `prepare_cleanup_merge`.
- `Sweepable` gains `merge_if_due()`; `flush()` now only seals.
`flush_pass` runs `merge_if_due` after a successful seal and reports its
outcome via the shared `report_merge` under the **cleanup** counters, so
dashboards attribute it correctly.
- The flush sweeper's three kind passes run with `tokio::join!` (the
cleanup sweeper has since #256): it now carries the generic merge, so a
slow rollout walk must not delay generic's tick.
- `POST /generic/{name}/merge-wal` uses the same split (the rollout
route already did).

No config or format changes.

## Verification

- `cargo clippy --workspace --all-targets -D warnings` clean; full
server suite passes (84).
- New tests in `sweeper::tests`:
  - count trigger merges at threshold and drains pending to 0;
  - no-op below threshold;
- `merge_after_generations = 0` is *disabled*, not threshold 0 (guards
the `threshold.max(1)` footgun) — the time trigger still drains;
- **a reader holding the shared lock for the entire merge does not block
it** — `merge_if_due` completes once the reader drops. Under the old
exclusive-lock-across-prepare shape this test deadlocks.
- Will go through the staging load test alongside #261 before merge is
requested.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Beinan Wang <>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
beinan added a commit that referenced this pull request Sep 23, 2026
…ed MergeWal task (#264)

## Problem

Generic stores were never registered with the master. Their MemWAL
merges ran only on each worker's own timers (`spawn_flush_sweeper` /
`spawn_global_sweeper`), so all N workers raced to commit their shard's
flushed generations into **one** base table.

In production (20 workers, ABFS) under storage-account throttling this
failed continuously:

```
Too many concurrent writers. Attempted 2 times, but failed on retry_timeout of 30.000 seconds.
  (lance/src/dataset/write/retry.rs)
```

Lance's commit-conflict retry re-runs the whole `merge_insert` and gives
up after 30 s wall clock. With 20 writers on one commit point and a
throttled object store, the second attempt never completes in time. The
generations that did not merge kept the read fan-out — every read opened
~500 generation datasets — which kept the account throttled. Serially
triggering `/merge-wal` by hand while worker timers were disabled still
failed on ~30% of shards.

The master's `MergeWal` task — per-target etcd lock, one task per store,
fan-out to `WORKER_ENDPOINTS` — exists to serialize exactly this, and
rollout stores have used it from the start. Generic stores just weren't
wired into it.

## Fix

Bring generic stores onto the existing master path; no new mechanism.

- **`MasterState`** opens `_registry.generic.lance` (same
`RolloutRegistry` type the data plane writes) alongside the rollout
registry; adds `generic_uri()` / `generic_store_options()`.
- **Stats scan** enumerates both registries. Generic stores are observed
(version, base row count, fragment count, pending WAL generations) via
`GenericStore::open_existing`, and their rows are written under a
**`generic:<name>`** target. Store names match
`[A-Za-z0-9_][A-Za-z0-9._-]*` so `:` is unambiguous, and carrying the
kind in the target string keeps the etcd task schema, dedupe keys and
per-target locks unchanged.
- **`sweep_merge_wal_inner`** therefore picks up generic rows over
`merge_wal_min_generations` with no change.
- **`run_merge_wal`** parses the target: `generic:` → `POST
/api/v1/generic/{name}/merge-wal`, bare →
`/api/v1/internal/merge-wal/{name}`.
- **Fan-out is now serial** instead of `join_all`. Every worker's merge
commits a new version of the same base table, so parallel fan-out was N
writers racing one commit point even with the task lock held. This
applies to rollout too.
- `Compact` / `IndexId` tasks and the compaction sweep remain
rollout-only (generic rows are skipped with a clear error if enqueued by
hand); retirement only considers rollout rows.

## Verification

- `cargo clippy --workspace --all-targets -D warnings` clean.
- Master suite: 42 default + 25 etcd-backed pass (`ETCD_TEST_ENDPOINTS`
against a local etcd 3.7), including four new tests:
- `merge_wal_routes_generic_targets_to_the_generic_endpoint` — a
`generic:gs` task hits `/api/v1/generic/gs/merge-wal` with the bare name
and never the rollout route;
- `merge_wal_calls_workers_serially` — three stub workers, shared
in-flight counter, asserts max in-flight is 1;
- `sweep_merge_wal_enqueues_over_threshold_and_dedupes` — extended: a
generic row over threshold is enqueued as `generic:ghot` alongside the
rollout row, and de-dupes on the second sweep;
  - `compaction_sweep_skips_generic_rows`;
- scanner `generic_store_is_observed_with_pending_wal` — 3 sealed adds →
row under `generic:g` with `pending_wal_generations == 3`, skipped on an
unchanged re-scan.

Deployment note: with the master driving generic merges, the worker-side
count trigger for generic should be off
(`ROLLOUT_MERGE_AFTER_GENERATIONS=0`) and the 300 s cleanup kept as a
fallback — the same posture rollout already runs with.

Refs #256, #257, #258 (supersedes the "unified maintainer" direction
there — the right answer was the existing master path), #263.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Beinan Wang <>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
beinan added a commit that referenced this pull request Sep 25, 2026
## 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](https://claude.com/claude-code)

Co-authored-by: Beinan Wang <>
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