Skip to content

feat(core): sample and warn on MemWAL pending generations per shard - #261

Merged
beinan merged 1 commit into
lance-format:mainfrom
beinan:fix/wal-pending-generations-observability
Sep 22, 2026
Merged

beinan merged 1 commit into
lance-format:mainfrom
beinan:fix/wal-pending-generations-observability

Conversation

@beinan

@beinan beinan commented Sep 22, 2026

Copy link
Copy Markdown
Collaborator

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

🤖 Generated with Claude Code

…ance-format#259)

Every LSM read (wal_shard_snapshots) now records each shard's
flushed_generations.len() into a rollout_wal_pending_generations
histogram and warns once per shard per read when it reaches
ROLLOUT_WAL_PENDING_WARN_GENERATIONS (default 256; 0 disables).

Every pending generation is a separate dataset the read path must open,
so this count is the read-amplification signal. One generic store
reached ~12.6k pending generations across 20 shards before anyone
noticed, by which point every read exhausted a 32 GiB worker; nothing in
logs or metrics pointed at the store until the crash logs were read by
hand. The warn carries the dataset URI and shard id on the span so it is
grep-able the way the oversized-blob warn is.

No behaviour change. The threshold is plumbed through every store kind's
options and the server config, and the metrics crate gives the histogram
count-scale buckets (8..4096) instead of the duration defaults.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@beinan
beinan merged commit 92377a9 into lance-format:main Sep 22, 2026
10 checks passed
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 26, 2026
…e-enqueueing them (#267)

Closes #266.

## Problem

A store whose base-table manifest names a fragment that no longer exists
fails every `MergeWal` and `Compact` at the same point, forever:

```
task failed target=<store> error="Not found: .../<store>.rollout.lance/data/<frag>.lance"
```

The scheduler has no memory of failure. `sweep_merge_wal_inner` /
`sweep_candidates_inner` re-enqueue every over-threshold store on every
tick (`enqueue` only de-dupes against queued/running), so a task that
just failed is back in the queue ten minutes later, and each attempt
holds a `TASK_CONCURRENCY` slot through a full serial worker fan-out
before failing. Production 2026-09-26: five such stores (one at
**7,985** pending generations) consumed roughly half of 6 masters' merge
slots; the merge queue for healthy stores went 0 → 109 in an hour.

## Fix

`TaskStore` records consecutive failures per `(kind, target)` in etcd
under `<prefix>/cooldown/<kind>/<target>`, leased for the cooldown
duration so it ages out on its own.

- After `TASK_COOLDOWN_AFTER_FAILURES` (default **3**), the sweeps skip
the target for `TASK_COOLDOWN_BASE_SECS` (default **600**), doubling per
further failure up to `TASK_COOLDOWN_MAX_SECS` (default **21600**). `0`
disables.
- A success clears the record. A manual `POST /tasks` is **not** gated —
operators can always force an attempt.
- Below the threshold the record only carries the count (leased for
`max`, so a slow trickle of unrelated failures doesn't accumulate
forever) and is not a cooldown.
- Entering cooldown emits `warn!` with the last error and
`master_task_cooldowns_total{kind}`; `GET /api/v1/scheduler/cooldowns`
returns `Vec<TaskCooldown>` (kind, target, failures, until_ms,
last_error) so a broken store is visible instead of silently eating
capacity.
- Bookkeeping is best-effort after the task's terminal state is
committed; a failure to write the cooldown key never fails the task
completion.

## Verification

- `cargo clippy --workspace --all-targets -D warnings` clean.
- etcd-backed suite 26/26 (`ETCD_TEST_ENDPOINTS` against local etcd
3.7), including new
`repeated_failures_cool_the_target_down_and_sweeps_skip_it`: threshold
2, no worker endpoints so every MergeWal fails; asserts the target is
*not* cooling after one failure, *is* after two, the sweep then enqueues
0, `list_cooldowns` reports it with `until_ms` set, and a manual enqueue
succeeds.
- Writing the test caught a bug in the first draft: a sub-threshold
record was treated as a cooldown. Fixed (`is_cooling_down` requires
`until_ms`).

Refs #264, #261.

🤖 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.

core: unbounded MemWAL read fan-out — warn and expose a metric when a shard's pending generations pile up

1 participant