fix(server): run the count-triggered WAL merge under the prepare/commit lock split - #263
Merged
Merged
Conversation
…it lock split (lance-format#257) 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, which reads up to MERGE_MAX_GENERATIONS generations from object storage before committing. The exclusive lock was held across that whole read (seconds to tens of seconds on ABFS), stalling every add/get/list on the store and the next 30s flush tick behind it. In production this surfaced as 'flush sweeper timed out' 52 times in 15 minutes and request handlers failing with 'Too many concurrent writers'. The cleanup path already used the prepare (read lock) / commit (write lock) split. This applies it to the count trigger too: - StorageBase::prepare_count_merge(): the shared-lock half of maybe_merge_own_shard, None when the trigger is disabled or not met. Passthroughs on RolloutStore and GenericStore. - Sweepable gains merge_if_due(); flush() now only seals. flush_pass runs merge_if_due after a successful seal and reports it under the cleanup counters, so a failing merge no longer shows up as a failing flush. - The flush sweeper's three kind passes run with tokio::join!, like the cleanup sweeper since lance-format#256, so the merge it now carries for generic is not queued behind a slow rollout walk. - POST /generic/{name}/merge-wal uses the same split; the rollout route already did. Tests: the count trigger merges at threshold, is a no-op below it, treats merge_after_generations=0 as disabled (not threshold 0), and completes while a reader holds the shared lock for the whole merge — which deadlocks under the old exclusive-lock-across-prepare shape. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This was referenced Sep 23, 2026
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>
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.
Closes #257.
Problem
Sweepable::flushfor rollout and generic sealed under the read lock, then took the write lock and calledmaybe_merge_own_shard()/maybe_merge_wal(). Those read up toMERGE_MAX_GENERATIONSgenerations from object storage before committing, so the exclusive lock was held across the whole read — seconds to tens of seconds on ABFS. Everyadd/get/liston the store stalled, and the next 30s flush tick queued behind it.Production after #256:
flush sweeper timed out52× in 15 min on one hot generic store, request handlers failing withToo 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 failedand counted underrollout_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 ofmaybe_merge_own_shard;Nonewhen the trigger is disabled (merge_after_generations == 0) or not met. Thin passthroughs onRolloutStoreandGenericStore, mirroringprepare_cleanup_merge.Sweepablegainsmerge_if_due();flush()now only seals.flush_passrunsmerge_if_dueafter a successful seal and reports its outcome via the sharedreport_mergeunder the cleanup counters, so dashboards attribute it correctly.tokio::join!(the cleanup sweeper has since fix(server): stop generic-store WAL merges starving behind rollout sweeps #256): it now carries the generic merge, so a slow rollout walk must not delay generic's tick.POST /generic/{name}/merge-waluses the same split (the rollout route already did).No config or format changes.
Verification
cargo clippy --workspace --all-targets -D warningsclean; full server suite passes (84).sweeper::tests:merge_after_generations = 0is disabled, not threshold 0 (guards thethreshold.max(1)footgun) — the time trigger still drains;merge_if_duecompletes once the reader drops. Under the old exclusive-lock-across-prepare shape this test deadlocks.🤖 Generated with Claude Code