From 1b0edd80e5706652c4f3dbab644f37a96ab51f23 Mon Sep 17 00:00:00 2001 From: Beinan Wang Date: Wed, 30 Sep 2026 08:07:27 +0000 Subject: [PATCH] fix: stop reads and merges from opening the whole shard at once 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 --- .../lance-context-core/src/datagen_store.rs | 3 + .../lance-context-core/src/generic_store.rs | 3 + crates/lance-context-core/src/lib.rs | 3 + .../lance-context-core/src/rollout_store.rs | 131 ++++++++++++++++++ crates/lance-context-core/src/store.rs | 5 + crates/lance-context-core/src/store_base.rs | 83 +++++++++++ crates/lance-context-master/src/scheduler.rs | 17 +++ crates/lance-context-master/src/state.rs | 3 + crates/lance-context-server/src/config.rs | 21 +++ crates/lance-context-server/src/error.rs | 12 ++ .../src/routes/datagen.rs | 1 + .../src/routes/generic.rs | 2 + .../src/routes/rollouts.rs | 2 + crates/lance-context-server/src/state.rs | 24 +++- crates/lance-context/src/unified_datagen.rs | 1 + crates/lance-context/src/unified_generic.rs | 2 + 16 files changed, 312 insertions(+), 1 deletion(-) diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index 1b83ea9..5e66edd 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -65,6 +65,8 @@ pub struct DatagenStoreOptions { /// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations` /// histogram is emitted regardless. pub pending_generations_warn: Option, + /// See `StorageBaseOptions::pending_generations_max`. + pub pending_generations_max: Option, /// Process-wide byte budget shared by every merge this process runs; a /// merge that cannot fit waits for another to release. `None` disables /// the bound. See [`crate::merge_budget`] for the design. @@ -114,6 +116,7 @@ impl DatagenStore { merge_max_generations: options.merge_max_generations, merge_max_bytes: options.merge_max_bytes, pending_generations_warn: options.pending_generations_warn, + pending_generations_max: options.pending_generations_max, merge_budget: options.merge_budget.clone(), session: None, schema: Arc::new(datagen_log_schema()), diff --git a/crates/lance-context-core/src/generic_store.rs b/crates/lance-context-core/src/generic_store.rs index 8520025..da36431 100644 --- a/crates/lance-context-core/src/generic_store.rs +++ b/crates/lance-context-core/src/generic_store.rs @@ -78,6 +78,8 @@ pub struct GenericStoreOptions { /// alarm. `None` uses the crate default (256); `Some(0)` disables the warn. /// The `rollout_wal_pending_generations` histogram is emitted regardless. pub pending_generations_warn: Option, + /// See `StorageBaseOptions::pending_generations_max`. + pub pending_generations_max: Option, /// Process-wide byte budget shared by every merge this process runs; a /// merge that cannot fit waits for another to release. `None` disables /// the bound. See [`crate::merge_budget`] for the design. @@ -186,6 +188,7 @@ impl GenericStore { merge_max_generations: options.merge_max_generations, merge_max_bytes: options.merge_max_bytes, pending_generations_warn: options.pending_generations_warn, + pending_generations_max: options.pending_generations_max, merge_budget: options.merge_budget.clone(), session: options.session, schema: create_schema, diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index a0aec7d..a46587b 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -21,6 +21,9 @@ pub mod serde; mod storage; mod store; mod store_base; +pub use store_base::{ + is_pending_generations_exceeded, DEFAULT_PENDING_GENERATIONS_MAX, PENDING_GENERATIONS_EXCEEDED, +}; // Request/DTO conversions, exported so the server does not keep its own copies. // These were duplicated verbatim between here and `routes/`; see #214. diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index 6857e46..97c782f 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -389,6 +389,8 @@ pub struct RolloutStoreOptions { /// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations` /// histogram is emitted regardless. pub pending_generations_warn: Option, + /// See `StorageBaseOptions::pending_generations_max`. + pub pending_generations_max: Option, /// Process-wide byte budget shared by every merge this process runs; a /// merge that cannot fit waits for another to release. `None` disables /// the bound. See [`crate::merge_budget`] for the design. @@ -480,6 +482,7 @@ impl RolloutStore { merge_max_generations, merge_max_bytes, pending_generations_warn, + pending_generations_max, merge_budget, session, } = options; @@ -492,6 +495,7 @@ impl RolloutStore { merge_max_generations, merge_max_bytes, pending_generations_warn, + pending_generations_max, merge_budget, session, schema: Arc::new(rollout_schema()), @@ -707,6 +711,11 @@ impl RolloutStore { self.base.create_key_btree_index().await } + /// See `StorageBase::extend_key_btree_index`. + pub async fn extend_id_btree_index(&mut self) -> LanceResult { + self.base.extend_key_btree_index().await + } + /// See `StorageBase::has_key_btree_index`. pub async fn has_id_btree_index(&self) -> LanceResult { self.base.has_key_btree_index().await @@ -2525,6 +2534,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, session: None, schema: legacy_schema.clone(), @@ -2846,6 +2856,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -2892,6 +2903,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }; @@ -2943,6 +2955,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -2996,6 +3009,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }; @@ -3219,6 +3233,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, ..Default::default() }, @@ -3283,6 +3298,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, ..Default::default() }, @@ -3385,6 +3401,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3445,6 +3462,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3486,6 +3504,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3623,6 +3642,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3711,6 +3731,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3799,6 +3820,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3848,6 +3870,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3919,6 +3942,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -3958,6 +3982,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -4018,6 +4043,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -4061,6 +4087,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -4272,6 +4299,107 @@ mod tests { }); } + /// Every merge appends fragments the id BTree does not cover, and + /// `merge_insert` full-scans exactly those. `extend_id_btree_index` + /// appends an index delta over them; a fully covered table is a no-op + /// and a table without the BTree is left alone. + #[test] + fn extend_id_btree_index_covers_fragments_added_since_build() { + use lance::index::DatasetIndexInternalExt as _; + let dir = TempDir::new().unwrap(); + let uri = dir.path().to_string_lossy().to_string(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let mut store = RolloutStore::open(&uri).await.unwrap(); + store.add(&[assistant_record("a-0")]).await.unwrap(); + store.flush().await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + + // No BTree yet: nothing to extend. + assert_eq!(store.extend_id_btree_index().await.unwrap(), 0); + + store.create_id_btree_index().await.unwrap(); + assert_eq!(store.extend_id_btree_index().await.unwrap(), 0); + + // Two merges land two fragments the index knows nothing about. + for id in ["a-1", "a-2"] { + store.add(&[assistant_record(id)]).await.unwrap(); + store.flush().await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + } + let unindexed = |s: &RolloutStore| { + let dataset = s.base.dataset.clone(); + async move { + dataset + .unindexed_fragments(ROLLOUT_ID_INDEX_NAME) + .await + .unwrap() + .len() + } + }; + assert_eq!(unindexed(&store).await, 2); + + assert_eq!(store.extend_id_btree_index().await.unwrap(), 2); + assert_eq!(unindexed(&store).await, 0); + assert!(store.has_id_btree_index().await.unwrap()); + assert_eq!(store.extend_id_btree_index().await.unwrap(), 0); + + // Rows are still merged correctly through the extended index. + store.add(&[assistant_record("a-0")]).await.unwrap(); + store.flush().await.unwrap(); + store.cleanup_own_shard().await.unwrap(); + assert_eq!(store.list(None, None).await.unwrap().len(), 3); + }); + } + + /// Reads open every flushed generation pending merge. Past the cap the + /// read is refused with a recognizable error (the server maps it to 503) + /// instead of holding a worker's memory hostage; at the cap it passes. + #[test] + fn reads_are_refused_past_the_pending_generation_cap() { + let dir = TempDir::new().unwrap(); + let uri = dir.path().to_string_lossy().to_string(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let mut store = RolloutStore::open_with_options( + &uri, + RolloutStoreOptions { + storage_options: None, + session: None, + shard_id: Some("rollout-0".to_string()), + merge_after_generations: None, + merge_max_generations: None, + merge_max_bytes: None, + pending_generations_warn: None, + pending_generations_max: Some(2), + merge_budget: None, + }, + ) + .await + .unwrap(); + + for id in ["a-0", "a-1"] { + store.add(&[assistant_record(id)]).await.unwrap(); + store.flush().await.unwrap(); + } + // At the cap: still readable. + assert_eq!(store.list(None, None).await.unwrap().len(), 2); + + store.add(&[assistant_record("a-2")]).await.unwrap(); + store.flush().await.unwrap(); + let err = store.list(None, None).await.expect_err("over the cap"); + assert!( + crate::store_base::is_pending_generations_exceeded(&err), + "{err}" + ); + assert!(err.to_string().contains("3 flushed generations"), "{err}"); + + // The merge drains the backlog and reads resume. + assert_eq!(store.cleanup_own_shard().await.unwrap(), 3); + assert_eq!(store.list(None, None).await.unwrap().len(), 3); + }); + } + #[test] fn create_id_btree_index_builds_and_is_idempotent() { // Building the BTree index on `id` must succeed even though the @@ -4638,6 +4766,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, ..Default::default() }, @@ -4738,6 +4867,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) @@ -4772,6 +4902,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, }, ) diff --git a/crates/lance-context-core/src/store.rs b/crates/lance-context-core/src/store.rs index 137f2a9..9acba05 100644 --- a/crates/lance-context-core/src/store.rs +++ b/crates/lance-context-core/src/store.rs @@ -329,6 +329,8 @@ pub struct ContextStoreOptions { /// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations` /// histogram is emitted regardless. pub pending_generations_warn: Option, + /// Refuse reads once more than this many flushed generations are pending merge. + pub pending_generations_max: Option, /// Process-wide byte budget shared by every merge this process runs; a /// merge that cannot fit waits for another to release. `None` disables /// the bound. See [`crate::merge_budget`] for the design. @@ -365,6 +367,7 @@ impl Default for ContextStoreOptions { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, // Read-your-write by default; see the field docs. seal_on_add: true, @@ -648,6 +651,7 @@ impl ContextStore { merge_max_generations: options.merge_max_generations, merge_max_bytes: options.merge_max_bytes, pending_generations_warn: options.pending_generations_warn, + pending_generations_max: options.pending_generations_max, merge_budget: options.merge_budget.clone(), session: None, schema: Arc::new(arrow_schema.clone()), @@ -2273,6 +2277,7 @@ impl ContextStore { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, // A compactor never appends, so the seal mode is irrelevant to it; // deferring keeps it from ever emitting a generation. diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index 7f67572..f5cf1be 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -161,6 +161,18 @@ async fn compact_files_incremental( /// every store, so all tables index their primary key identically. pub(crate) const ID_INDEX_NAME: &str = "id_idx"; +/// Error-message prefix for a read refused because a shard has too many +/// flushed generations pending merge. Servers match on it to answer 503. +pub const PENDING_GENERATIONS_EXCEEDED: &str = "too many pending WAL generations"; + +/// Default for [`StorageBaseOptions::pending_generations_max`]. +pub const DEFAULT_PENDING_GENERATIONS_MAX: usize = 4096; + +/// Whether `err` is a read refused by the pending-generations cap. +pub fn is_pending_generations_exceeded(err: &LanceError) -> bool { + err.to_string().contains(PENDING_GENERATIONS_EXCEEDED) +} + /// What a [`StorageBase::flush`] actually did, for the `outcome` metric label. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum FlushOutcome { @@ -245,6 +257,11 @@ pub(crate) struct StorageBaseOptions { /// pending merge (sampled on every LSM read). `None` uses the crate default /// (256); `Some(0)` disables the warning. The metric is always emitted. pub pending_generations_warn: Option, + /// Refuse a read whose shard has more flushed generations pending than + /// this, instead of opening every one of them. `None` uses the crate + /// default (4096); `Some(0)` disables the cap. See + /// [`StorageBase::wal_shard_snapshots`]. + pub pending_generations_max: Option, /// Process-wide byte budget shared by every merge this process runs. /// `merge_max_bytes` bounds one merge; this bounds all of them together, /// and a merge that cannot fit waits for another to release. `None` @@ -313,6 +330,9 @@ pub(crate) struct StorageBase { merge_max_bytes: usize, /// Per-shard pending-generation count at which reads warn; `0` disables. pending_generations_warn: usize, + /// Refuse reads that would open more flushed generations than this per + /// shard (0 = unbounded). See [`Self::wal_shard_snapshots`]. + pending_generations_max: usize, /// Process-wide merge byte budget; `None` means unbounded. merge_budget: Option>, /// Timestamp of the last successful [`Self::compact`] on this handle. @@ -372,6 +392,7 @@ impl StorageBase { merge_max_generations, merge_max_bytes, pending_generations_warn, + pending_generations_max, merge_budget, session, schema, @@ -404,6 +425,7 @@ impl StorageBase { merge_max_generations, merge_max_bytes, pending_generations_warn, + pending_generations_max, merge_budget, session, schema, @@ -430,6 +452,7 @@ impl StorageBase { merge_max_generations, merge_max_bytes, pending_generations_warn, + pending_generations_max, merge_budget, session, schema, @@ -458,6 +481,8 @@ impl StorageBase { merge_max_bytes: merge_max_bytes.unwrap_or(DEFAULT_MERGE_MAX_BYTES), pending_generations_warn: pending_generations_warn .unwrap_or(DEFAULT_PENDING_GENERATIONS_WARN), + pending_generations_max: pending_generations_max + .unwrap_or(DEFAULT_PENDING_GENERATIONS_MAX), merge_budget, last_compaction: None, total_compactions: 0, @@ -1477,6 +1502,48 @@ impl StorageBase { self.reload().await } + /// Extend the key column's BTree index over every base-table fragment it + /// does not yet cover, appending an index delta rather than rebuilding. + /// + /// Lance's `merge_insert` probes the index for the fragments it covers and + /// **scans** every fragment it does not (its plan is a `Union` of the two). + /// Every WAL merge appends fragments the index has never seen, so between + /// rebuilds each merge reads those fragments whole: 6.5 GB for a 2,046-row + /// merge into a 22-fragment table was observed, and that read is what + /// spiked workers to 20-29 GiB after the BTree itself was in place. + /// Calling this before a merge keeps the unindexed set empty, so the scan + /// arm reads nothing. + /// + /// Returns how many fragments were unindexed beforehand (0 = no commit). + /// A table without the BTree at all is left alone; the caller builds it + /// with [`Self::create_key_btree_index`] first. + pub async fn extend_key_btree_index(&mut self) -> LanceResult { + use lance::index::{DatasetIndexExt as _, DatasetIndexInternalExt as _}; + use lance_index::optimize::OptimizeOptions; + + self.ensure_writable()?; + self.reload().await?; + if !self.has_key_btree_index().await? { + return Ok(0); + } + let unindexed = self.dataset.unindexed_fragments(ID_INDEX_NAME).await?.len(); + if unindexed == 0 { + return Ok(0); + } + self.dataset + .optimize_indices( + &OptimizeOptions::append().index_names(vec![ID_INDEX_NAME.to_string()]), + ) + .await?; + self.reload().await?; + info!( + column = %self.key_column, + fragments = unindexed, + "extended key index over newly appended fragments" + ); + Ok(unindexed) + } + /// Whether the base table's key column has a BTree index (the one /// `merge_insert` can use). A ZoneMap under the same name does not count. pub async fn has_key_btree_index(&self) -> LanceResult { @@ -1613,6 +1680,7 @@ impl StorageBase { let branch_path = self.dataset.branch_location().path.clone(); let shard_ids = self.dataset.list_mem_wal_latest_shard_ids().await?; let warn_at = self.pending_generations_warn; + let max_at = self.pending_generations_max; let uri: Arc = Arc::from(self.dataset.uri()); let snapshots: Vec> = stream::iter(shard_ids) @@ -1632,6 +1700,21 @@ impl StorageBase { }; let pending = manifest.flushed_generations.len(); observe_value!(crate::metrics::ROLLOUT_WAL_PENDING_GENERATIONS, pending); + if max_at != 0 && pending > max_at { + // A read over this shard would open `pending` datasets + // and hold their metadata at once. Past the cap that is + // gigabytes per request (a 16k-generation shard took a + // worker to its 32 GiB limit on one point lookup), and + // the only cure is the merge that is already behind. + // Fail the read fast so the caller retries after the + // merge instead of taking the worker down with it. + metrics::counter!("rollout_reads_refused_pending_total").increment(1); + return Err(LanceError::io(format!( + "{PENDING_GENERATIONS_EXCEEDED}: shard {shard_id} has {pending} \ + flushed generations pending merge (cap {max_at}); retry after \ + the merge catches up" + ))); + } if warn_at != 0 && pending >= warn_at { // Every pending generation is one more dataset this read // must open; past this point the shard's owner has stopped diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 73fd151..c9c028f 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -396,6 +396,23 @@ async fn ensure_id_btree_index(state: &Arc, name: &str) -> Result<( .await .map_err(|e| e.to_string())? { + // The BTree exists but every merge and compaction since it was built + // added fragments it does not cover, and `merge_insert` full-scans + // those. Append an index delta over them so the probe stays a probe. + let started = std::time::Instant::now(); + let covered = store + .extend_id_btree_index() + .await + .map_err(|e| format!("extending id index before merge: {e}"))?; + if covered > 0 { + metrics::counter!("master_merge_wal_index_extended_total").increment(1); + tracing::info!( + target = %name, + fragments = covered, + elapsed_secs = started.elapsed().as_secs(), + "extended id BTree index over fragments added since it was built" + ); + } return Ok(()); } let started = std::time::Instant::now(); diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index ede5157..425d85d 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -191,6 +191,9 @@ impl MasterState { pub(crate) fn rollout_store_options(&self) -> RolloutStoreOptions { RolloutStoreOptions { session: self.rollout_session.clone(), + // The master observes and merges shards *because* they are behind; + // the worker-side read cap must never hide those from it. + pending_generations_max: Some(0), ..Default::default() } } diff --git a/crates/lance-context-server/src/config.rs b/crates/lance-context-server/src/config.rs index 3185e58..5ef4fd9 100644 --- a/crates/lance-context-server/src/config.rs +++ b/crates/lance-context-server/src/config.rs @@ -62,6 +62,18 @@ pub struct ServerConfig { )] pub rollout_wal_pending_warn_generations: usize, + /// Refuse a read whose MemWAL shard has more than this many flushed + /// generations pending merge (HTTP 503) instead of + /// opening every one of them. A read over a 16k-generation shard held a + /// worker at its 32 GiB memory limit; the merge that fixes it is the same + /// merge the read was competing with. `0` disables the cap. + #[arg( + long, + env = "ROLLOUT_WAL_PENDING_MAX_GENERATIONS", + default_value = "4096" + )] + pub rollout_wal_pending_max_generations: usize, + /// Process-wide byte budget for MemWAL merges, shared by every merge /// this worker runs regardless of what triggered it (its own sweepers, the /// count trigger, the manual route, or the master's fan-out). @@ -72,6 +84,15 @@ pub struct ServerConfig { #[arg(long, env = "ROLLOUT_MERGE_MEMORY_BYTES", default_value = "3221225472")] pub rollout_merge_memory_bytes: usize, + /// Maximum merge-wal requests this worker runs at once; later ones wait + /// for a slot. The master's `MERGE_WAL_CONCURRENCY` bounds tasks *per + /// master*, and every task fans out to every worker, so six masters at 4 + /// put 24 merges on each worker at the same time -- each with its own + /// `LANCE_MEM_POOL_SIZE x LANCE_CPU_THREADS` execution pool, which the + /// merge memory budget does not count. Default 4; `0` disables. + #[arg(long, env = "ROLLOUT_MERGE_CONCURRENCY", default_value_t = 4)] + pub rollout_merge_concurrency: usize, + /// Interval, in seconds, for the periodic per-shard WAL cleanup task. When /// non-zero, the global sweeper folds this instance's flushed MemWAL /// generations into the base table on a schedule — the *time* half of the diff --git a/crates/lance-context-server/src/error.rs b/crates/lance-context-server/src/error.rs index 6e48cbb..2cd4809 100644 --- a/crates/lance-context-server/src/error.rs +++ b/crates/lance-context-server/src/error.rs @@ -30,6 +30,9 @@ impl AppError { // Checked before the typed match: the compaction-in-progress signal is an // Arrow-wrapped error with no dedicated variant, so only its text // distinguishes it. + if lance_context_core::is_pending_generations_exceeded(&err) { + return AppError::Overloaded(err.to_string()); + } if err.to_string().contains("already in progress") { return AppError::CompactionInProgress; } @@ -109,6 +112,15 @@ mod tests { )); } + #[test] + fn pending_generation_cap_is_overloaded() { + let err = LanceError::io(format!( + "{}: shard rollout-0 has 5000 flushed generations pending merge (cap 4096)", + lance_context_core::PENDING_GENERATIONS_EXCEEDED + )); + assert!(matches!(AppError::from_lance(err), AppError::Overloaded(_))); + } + #[test] fn compaction_in_progress_is_detected_from_arrow_variant() { // Reproduce how the core raises it: an ArrowError folded into Lance's diff --git a/crates/lance-context-server/src/routes/datagen.rs b/crates/lance-context-server/src/routes/datagen.rs index 6268a76..8974613 100644 --- a/crates/lance-context-server/src/routes/datagen.rs +++ b/crates/lance-context-server/src/routes/datagen.rs @@ -49,6 +49,7 @@ pub async fn create_datagen_store( merge_max_generations: Some(state.rollout_merge_max_generations), merge_max_bytes: Some(state.rollout_merge_max_bytes), pending_generations_warn: Some(state.rollout_wal_pending_warn_generations), + pending_generations_max: Some(state.rollout_wal_pending_max_generations), merge_budget: state.merge_budget.clone(), cleanup_interval_secs: None, }; diff --git a/crates/lance-context-server/src/routes/generic.rs b/crates/lance-context-server/src/routes/generic.rs index bf48706..c5fe858 100644 --- a/crates/lance-context-server/src/routes/generic.rs +++ b/crates/lance-context-server/src/routes/generic.rs @@ -60,6 +60,7 @@ pub async fn create_generic_store( merge_max_generations: Some(state.rollout_merge_max_generations), merge_max_bytes: Some(state.rollout_merge_max_bytes), pending_generations_warn: Some(state.rollout_wal_pending_warn_generations), + pending_generations_max: Some(state.rollout_wal_pending_max_generations), merge_budget: state.merge_budget.clone(), session: None, seal_on_add: req.seal_on_add, @@ -272,6 +273,7 @@ pub async fn merge_generic_wal( State(state): State>, Path(name): Path, ) -> Result, AppError> { + let _slot = state.acquire_merge_slot().await; let store = state.get_or_open_generic_store(&name).await?; // Same prepare/commit split as the sweeper: the object-storage read of the // generations runs under the shared lock so the store keeps serving. diff --git a/crates/lance-context-server/src/routes/rollouts.rs b/crates/lance-context-server/src/routes/rollouts.rs index de75dd1..5f401af 100644 --- a/crates/lance-context-server/src/routes/rollouts.rs +++ b/crates/lance-context-server/src/routes/rollouts.rs @@ -200,6 +200,7 @@ pub async fn create_rollout_store( merge_max_generations: Some(state.rollout_merge_max_generations), merge_max_bytes: Some(state.rollout_merge_max_bytes), pending_generations_warn: Some(state.rollout_wal_pending_warn_generations), + pending_generations_max: Some(state.rollout_wal_pending_max_generations), merge_budget: state.merge_budget.clone(), session: state.rollout_session.clone(), }; @@ -727,6 +728,7 @@ pub async fn merge_wal( State(state): State>, Path(name): Path, ) -> Result, AppError> { + let _slot = state.acquire_merge_slot().await; let store_lock = state.get_or_open_rollout_store(&name).await?; // Split by lock scope: seal + read every flushed generation under the // *read* lock so ingest on this store keeps running, then take the write diff --git a/crates/lance-context-server/src/state.rs b/crates/lance-context-server/src/state.rs index 5affabc..a34a88d 100644 --- a/crates/lance-context-server/src/state.rs +++ b/crates/lance-context-server/src/state.rs @@ -11,7 +11,7 @@ use lance_context_core::{ RolloutStore, RolloutStoreOptions, Session, }; use lru::LruCache; -use tokio::sync::{Mutex, OwnedMutexGuard, RwLock}; +use tokio::sync::{Mutex, OwnedMutexGuard, RwLock, Semaphore}; use tokio::task::JoinHandle; use crate::config::ServerConfig; @@ -123,9 +123,12 @@ pub struct AppState { pub rollout_merge_max_bytes: usize, /// Per-shard pending-generation count at which reads warn; `0` disables. pub rollout_wal_pending_warn_generations: usize, + pub rollout_wal_pending_max_generations: usize, /// Process-wide merge memory budget shared by every store; `None` when /// disabled. See `lance_context_core::merge_budget`. pub merge_budget: Option>, + /// Worker-wide cap on concurrent merge-wal requests; `None` when disabled. + pub merge_slots: Option>, /// Periodic per-shard WAL-cleanup interval in seconds; `0` disables the /// global sweeper. See [`Self::spawn_global_sweeper`]. pub rollout_cleanup_interval_secs: u64, @@ -285,6 +288,17 @@ fn build_rollout_session(cache_bytes: usize) -> Option> { } impl AppState { + /// Wait for a merge slot (see `ServerConfig::rollout_merge_concurrency`). + pub async fn acquire_merge_slot(&self) -> Option { + let slots = self.merge_slots.as_ref()?; + let wait = std::time::Instant::now(); + // Only fails if the semaphore is closed, which we never do. + let permit = Arc::clone(slots).acquire_owned().await.ok()?; + metrics::histogram!("rollout_wal_merge_slot_wait_seconds") + .record(wait.elapsed().as_secs_f64()); + Some(permit) + } + /// Build the shared server state, opening (or creating) the rollout registry /// under `data_dir`. Async because opening the registry touches storage. pub async fn new(config: ServerConfig) -> Result { @@ -318,8 +332,11 @@ impl AppState { rollout_merge_max_generations: config.rollout_merge_max_generations, rollout_merge_max_bytes: config.rollout_merge_max_bytes, rollout_wal_pending_warn_generations: config.rollout_wal_pending_warn_generations, + rollout_wal_pending_max_generations: config.rollout_wal_pending_max_generations, merge_budget: (config.rollout_merge_memory_bytes > 0) .then(|| MergeMemoryBudget::new(config.rollout_merge_memory_bytes)), + merge_slots: (config.rollout_merge_concurrency > 0) + .then(|| Arc::new(Semaphore::new(config.rollout_merge_concurrency))), rollout_cleanup_interval_secs: config.rollout_cleanup_interval_secs, rollout_flush_interval_secs: config.rollout_flush_interval_secs, blob_budget, @@ -386,7 +403,9 @@ impl AppState { rollout_merge_max_generations: 8, rollout_merge_max_bytes: 1024 * 1024 * 1024, rollout_wal_pending_warn_generations: 256, + rollout_wal_pending_max_generations: 0, merge_budget: None, + merge_slots: None, rollout_cleanup_interval_secs: 0, rollout_flush_interval_secs: 0, blob_budget: None, @@ -423,6 +442,7 @@ impl AppState { merge_max_generations: Some(self.rollout_merge_max_generations), merge_max_bytes: Some(self.rollout_merge_max_bytes), pending_generations_warn: Some(self.rollout_wal_pending_warn_generations), + pending_generations_max: Some(self.rollout_wal_pending_max_generations), merge_budget: self.merge_budget.clone(), session: self.rollout_session.clone(), } @@ -613,6 +633,7 @@ impl AppState { merge_max_generations: Some(self.rollout_merge_max_generations), merge_max_bytes: Some(self.rollout_merge_max_bytes), pending_generations_warn: Some(self.rollout_wal_pending_warn_generations), + pending_generations_max: Some(self.rollout_wal_pending_max_generations), merge_budget: self.merge_budget.clone(), cleanup_interval_secs: None, } @@ -757,6 +778,7 @@ impl AppState { merge_max_generations: Some(self.rollout_merge_max_generations), merge_max_bytes: Some(self.rollout_merge_max_bytes), pending_generations_warn: Some(self.rollout_wal_pending_warn_generations), + pending_generations_max: Some(self.rollout_wal_pending_max_generations), merge_budget: self.merge_budget.clone(), session: self.rollout_session.clone(), seal_on_add, diff --git a/crates/lance-context/src/unified_datagen.rs b/crates/lance-context/src/unified_datagen.rs index c49b199..6fd6f31 100644 --- a/crates/lance-context/src/unified_datagen.rs +++ b/crates/lance-context/src/unified_datagen.rs @@ -39,6 +39,7 @@ impl DatagenStore { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, cleanup_interval_secs: None, }; diff --git a/crates/lance-context/src/unified_generic.rs b/crates/lance-context/src/unified_generic.rs index c448618..3af2b5a 100644 --- a/crates/lance-context/src/unified_generic.rs +++ b/crates/lance-context/src/unified_generic.rs @@ -42,6 +42,7 @@ impl GenericStore { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, session: None, seal_on_add, @@ -65,6 +66,7 @@ impl GenericStore { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + pending_generations_max: None, merge_budget: None, session: None, seal_on_add,