From c25b54c5f436e7471872b24eb8e728ae6c68f60a Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Fri, 25 Sep 2026 20:35:51 +0000 Subject: [PATCH] feat(core): process-wide byte budget for MemWAL merges 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 (#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 --- .../lance-context-core/src/datagen_store.rs | 6 + .../lance-context-core/src/generic_store.rs | 87 +++++++ crates/lance-context-core/src/lib.rs | 2 + crates/lance-context-core/src/merge_budget.rs | 239 ++++++++++++++++++ crates/lance-context-core/src/metrics.rs | 6 + .../lance-context-core/src/rollout_store.rs | 27 ++ crates/lance-context-core/src/store.rs | 8 + crates/lance-context-core/src/store_base.rs | 68 ++++- crates/lance-context-metrics/src/lib.rs | 1 + crates/lance-context-server/src/config.rs | 11 + .../src/routes/datagen.rs | 1 + .../src/routes/generic.rs | 1 + .../src/routes/rollouts.rs | 1 + crates/lance-context-server/src/state.rs | 13 +- crates/lance-context/src/unified_datagen.rs | 1 + crates/lance-context/src/unified_generic.rs | 2 + 16 files changed, 469 insertions(+), 5 deletions(-) create mode 100644 crates/lance-context-core/src/merge_budget.rs diff --git a/crates/lance-context-core/src/datagen_store.rs b/crates/lance-context-core/src/datagen_store.rs index aa46a468..69b49704 100644 --- a/crates/lance-context-core/src/datagen_store.rs +++ b/crates/lance-context-core/src/datagen_store.rs @@ -35,6 +35,7 @@ use crate::datagen::{ DatagenRootItemStatuses, DatagenRunOverview, DatagenStepCursor, DatagenStepKind, DatagenStreamWriter, DatagenValue, DatagenWriteContext, FoldedDatagenItem, }; +use crate::merge_budget::MergeMemoryBudget; use crate::store::{ column_as, column_as_optional, timestamp_from_micros, CompactionConfig, CompactionStats, }; @@ -64,6 +65,10 @@ pub struct DatagenStoreOptions { /// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations` /// histogram is emitted regardless. pub pending_generations_warn: 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. + pub merge_budget: Option>, /// Periodically merge this writer's pending generations. `None` or zero /// disables the timer. pub cleanup_interval_secs: Option, @@ -109,6 +114,7 @@ impl DatagenStore { merge_max_generations: options.merge_max_generations, merge_max_bytes: options.merge_max_bytes, pending_generations_warn: options.pending_generations_warn, + merge_budget: options.merge_budget.clone(), session: None, schema: Arc::new(datagen_log_schema()), // Datagen keys on `event_id`, not `id`: event ids are derived diff --git a/crates/lance-context-core/src/generic_store.rs b/crates/lance-context-core/src/generic_store.rs index 00b97eef..6cf3fd90 100644 --- a/crates/lance-context-core/src/generic_store.rs +++ b/crates/lance-context-core/src/generic_store.rs @@ -32,6 +32,7 @@ use lance::{Error as LanceError, Result as LanceResult}; use lance_index::mem_wal::ShardManifest; use crate::generic_codec::{batch_to_rows, rows_to_batch, Row}; +use crate::merge_budget::MergeMemoryBudget; use crate::store::{CompactionConfig, CompactionStats}; use crate::store_base::{ListSource, PreparedMerge, StorageBase, StorageBaseOptions}; use lance_context_api::schema_spec::{SchemaSpec, ID_COLUMN}; @@ -77,6 +78,10 @@ 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, + /// 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. + pub merge_budget: Option>, /// Shared, capacity-bounded Lance session. pub session: Option>, /// Whether [`GenericStore::add`] seals before returning, making the rows it @@ -181,6 +186,7 @@ impl GenericStore { merge_max_generations: options.merge_max_generations, merge_max_bytes: options.merge_max_bytes, pending_generations_warn: options.pending_generations_warn, + merge_budget: options.merge_budget.clone(), session: options.session, schema: create_schema, // Always `id`: the LSM merge key, which `SchemaSpec::validate` @@ -786,6 +792,87 @@ mod tests { }); } + /// Two stores share one process-wide merge budget sized so that only one + /// merge's initial reservation fits. The second merge must wait until the + /// first commits and releases, then complete. This is the bound that was + /// missing when 17 of 20 production workers OOMKilled at once: each merge + /// was capped, their sum was not. + #[test] + fn concurrent_merges_queue_on_the_shared_memory_budget() { + use crate::merge_budget::MergeMemoryBudget; + use std::time::Duration; + use tokio::sync::oneshot; + + let dir = TempDir::new().unwrap(); + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on(async { + // 1 MiB budget, and each merge asks for min(merge_max_bytes, budget) + // up front, so the budget admits exactly one merge at a time. + let budget = MergeMemoryBudget::new(1024 * 1024); + let opts = |shard: &str| GenericStoreOptions { + seal_on_add: true, + shard_id: Some(shard.to_string()), + merge_max_bytes: Some(1024 * 1024), + merge_budget: Some(budget.clone()), + ..Default::default() + }; + let uri_a = dir.path().join("a").to_string_lossy().to_string(); + let uri_b = dir.path().join("b").to_string_lossy().to_string(); + let a = GenericStore::open(&uri_a, spec(), opts("a")).await.unwrap(); + let b = GenericStore::open(&uri_b, spec(), opts("b")).await.unwrap(); + for i in 0..3 { + a.add(&[row(json!({"id": format!("a{i}")}))]).await.unwrap(); + b.add(&[row(json!({"id": format!("b{i}")}))]).await.unwrap(); + } + + // Prepare A's merge: it takes the whole budget and holds it inside + // the PreparedMerge until we commit. + let prepared_a = a + .prepare_cleanup_merge() + .await + .unwrap() + .expect("A has pending"); + assert_eq!(budget.reserved(), 1024 * 1024); + + // B's prepare must block on the budget. + let (started_tx, started_rx) = oneshot::channel(); + let b_task = tokio::spawn(async move { + started_tx.send(()).unwrap(); + let prepared = b + .prepare_cleanup_merge() + .await + .unwrap() + .expect("B has pending"); + (b, prepared) + }); + started_rx.await.unwrap(); + tokio::time::sleep(Duration::from_millis(300)).await; + assert!( + !b_task.is_finished(), + "B's merge must wait while A holds the whole budget" + ); + + // Commit A: consumes the batches and releases the reservation. + let mut a = a; + let (ms, m, p) = prepared_a; + a.commit_prepared_merge(&ms, &m, p).await.unwrap(); + assert_eq!(a.pending_wal_generations().await.unwrap(), 0); + + // B proceeds now. + let (mut b, (ms, m, p)) = tokio::time::timeout(Duration::from_secs(30), b_task) + .await + .expect("B must proceed once A releases its reservation") + .unwrap(); + b.commit_prepared_merge(&ms, &m, p).await.unwrap(); + assert_eq!(b.pending_wal_generations().await.unwrap(), 0); + assert_eq!( + budget.reserved(), + 0, + "everything released after both commits" + ); + }); + } + #[test] fn wal_generations_merge_into_the_base_table() { let dir = TempDir::new().unwrap(); diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 586fae01..279b969c 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -10,6 +10,7 @@ mod export; pub mod generic_codec; mod generic_store; mod id; +pub mod merge_budget; pub mod metrics; mod namespace; mod record; @@ -52,6 +53,7 @@ pub use export::{ pub use generic_codec::{batch_to_rows, ids_from_batch, rows_to_batch, Row}; pub use generic_store::{GenericStore, GenericStoreOptions}; pub use id::{generate_id, new_uuid}; +pub use merge_budget::{MergeMemoryBudget, MergeReservation}; pub use namespace::{ContextNamespace, PartitionInfo, PartitionSelector, PartitionSpec}; pub use record::{ ContextRecord, LifecycleQueryOptions, MetadataFilter, RecordFilters, RecordPatch, Relationship, diff --git a/crates/lance-context-core/src/merge_budget.rs b/crates/lance-context-core/src/merge_budget.rs new file mode 100644 index 00000000..15d73d8f --- /dev/null +++ b/crates/lance-context-core/src/merge_budget.rs @@ -0,0 +1,239 @@ +//! Process-wide byte budget for MemWAL merges. +//! +//! A merge reads up to `merge_max_bytes` of Arrow batches into memory before it +//! commits. That bounds one merge. 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. With the master scheduling merges +//! for every store 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 of each other. +//! +//! [`MergeMemoryBudget`] is the missing bound. It is measured in **bytes**, not +//! merges: twenty tiny merges and one huge merge are priced by what they hold. +//! A merge reserves from it before reading and grows its reservation as it +//! reads; when the budget is exhausted the next merge **waits** for a release +//! rather than being rejected, so a master fanning out to a busy worker sees a +//! slower worker, not a failure. The reservation is RAII and travels inside +//! [`crate::store_base::PreparedMerge`], so it is released exactly when the +//! merged batches are dropped after commit. +//! +//! # Why it cannot deadlock +//! +//! Every merge takes its full initial reservation -- +//! `min(merge_max_bytes, budget)` -- **before** it reads anything, in one +//! atomic acquire. A merge therefore never holds part of what it needs while +//! waiting for the rest. The only case that grows a reservation mid-read is a +//! single generation larger than `merge_max_bytes`, which the merge folds whole +//! (a generation is indivisible). At that point it already holds at least as +//! much as any other merge could be waiting for, so growing cannot form a +//! cycle. If it needs more than the whole budget it is admitted anyway once it +//! is the sole holder, mirroring the "lone oversized request" rule of the blob +//! budget: an oversized generation must make progress or the shard wedges. + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; + +use tokio::sync::Semaphore; + +/// Reservation granularity. One permit is one MiB; the semaphore's permit +/// count is bounded (`Semaphore::MAX_PERMITS`), and byte-granular permits would +/// overflow it for budgets past a few GiB. +const PERMIT_BYTES: usize = 1024 * 1024; + +/// Byte budget shared by every merge in a process. +#[derive(Debug)] +pub struct MergeMemoryBudget { + permits: Semaphore, + limit: usize, + /// Bytes currently reserved, for the gauge. Tracked separately because + /// `Semaphore::available_permits` rounds to MiB. + reserved: AtomicUsize, +} + +/// RAII reservation held by an in-flight merge. Releases on drop. +#[derive(Debug)] +pub struct MergeReservation { + budget: Arc, + permits: u32, + bytes: usize, +} + +impl MergeMemoryBudget { + /// A budget admitting up to `limit` bytes of concurrently buffered merge + /// data. `limit` is rounded up to whole MiB. + #[must_use] + pub fn new(limit: usize) -> Arc { + let permits = Self::permits_for(limit).max(1); + Arc::new(Self { + permits: Semaphore::new(permits as usize), + limit: permits as usize * PERMIT_BYTES, + reserved: AtomicUsize::new(0), + }) + } + + /// The configured limit, in bytes (rounded up to whole MiB). + #[must_use] + pub fn limit(&self) -> usize { + self.limit + } + + /// Bytes currently reserved across all in-flight merges. + #[must_use] + pub fn reserved(&self) -> usize { + self.reserved.load(Ordering::Acquire) + } + + fn permits_for(bytes: usize) -> u32 { + let permits = bytes.div_ceil(PERMIT_BYTES); + u32::try_from(permits).unwrap_or(u32::MAX) + } + + /// Reserve `bytes`, waiting until they fit. A request larger than the whole + /// budget waits for the budget to be idle and then takes all of it: an + /// oversized merge must still make progress, so it is admitted alone. + pub async fn reserve(self: &Arc, bytes: usize) -> MergeReservation { + let want = Self::permits_for(bytes); + let total = self.permits_total(); + let permits = want.min(total); + // `acquire_many` only fails if the semaphore is closed, which we never do. + let permit = self + .permits + .acquire_many(permits) + .await + .expect("merge memory budget semaphore is never closed"); + permit.forget(); + self.reserved.fetch_add(bytes, Ordering::AcqRel); + MergeReservation { + budget: Arc::clone(self), + permits, + bytes, + } + } + + fn permits_total(&self) -> u32 { + Self::permits_for(self.limit) + } +} + +impl MergeReservation { + /// Bytes this reservation currently accounts for. + #[must_use] + pub fn bytes(&self) -> usize { + self.bytes + } + + /// Grow the reservation to cover `new_total` bytes, waiting for the + /// additional permits if the budget is exhausted. A no-op when `new_total` + /// does not exceed what is already held. The extra is capped at the budget: + /// a reservation never needs more permits than exist, and a holder that + /// already has them all is the sole merge in flight. + pub async fn grow_to(&mut self, new_total: usize) { + if new_total <= self.bytes { + return; + } + let want_total = MergeMemoryBudget::permits_for(new_total).min(self.budget.permits_total()); + let extra = want_total.saturating_sub(self.permits); + if extra > 0 { + let permit = self + .budget + .permits + .acquire_many(extra) + .await + .expect("merge memory budget semaphore is never closed"); + permit.forget(); + self.permits += extra; + } + self.budget + .reserved + .fetch_add(new_total - self.bytes, Ordering::AcqRel); + self.bytes = new_total; + } +} + +impl Drop for MergeReservation { + fn drop(&mut self) { + self.budget.permits.add_permits(self.permits as usize); + self.budget.reserved.fetch_sub(self.bytes, Ordering::AcqRel); + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::time::Duration; + + #[tokio::test] + async fn reservations_release_on_drop() { + let budget = MergeMemoryBudget::new(4 * PERMIT_BYTES); + let a = budget.reserve(2 * PERMIT_BYTES).await; + assert_eq!(budget.reserved(), 2 * PERMIT_BYTES); + drop(a); + assert_eq!(budget.reserved(), 0); + assert_eq!(budget.permits.available_permits(), 4); + } + + /// The property that matters in production: with the budget full, the next + /// merge waits instead of proceeding, and resumes when a holder releases. + #[tokio::test] + async fn exhausted_budget_makes_the_next_merge_wait() { + let budget = MergeMemoryBudget::new(2 * PERMIT_BYTES); + let held = budget.reserve(2 * PERMIT_BYTES).await; + + let waiter = { + let budget = budget.clone(); + tokio::spawn(async move { budget.reserve(PERMIT_BYTES).await }) + }; + tokio::time::sleep(Duration::from_millis(100)).await; + assert!(!waiter.is_finished(), "must wait while the budget is full"); + + drop(held); + let got = tokio::time::timeout(Duration::from_secs(5), waiter) + .await + .expect("waiter must proceed once bytes are released") + .unwrap(); + assert_eq!(got.bytes(), PERMIT_BYTES); + } + + /// A single merge larger than the whole budget is admitted (alone) rather + /// than waiting forever; an oversized generation must be folded whole. + #[tokio::test] + async fn oversized_request_takes_the_whole_budget_and_proceeds() { + let budget = MergeMemoryBudget::new(2 * PERMIT_BYTES); + let big = tokio::time::timeout(Duration::from_secs(2), budget.reserve(10 * PERMIT_BYTES)) + .await + .expect("an oversized request must not wait forever on an idle budget"); + assert_eq!(budget.permits.available_permits(), 0); + drop(big); + assert_eq!(budget.permits.available_permits(), 2); + } + + #[tokio::test] + async fn grow_to_waits_for_the_extra_and_is_idempotent_downward() { + let budget = MergeMemoryBudget::new(3 * PERMIT_BYTES); + let mut a = budget.reserve(PERMIT_BYTES).await; + let b = budget.reserve(2 * PERMIT_BYTES).await; + + let grow = { + // Growing `a` to 2 MiB needs one more permit; none free until `b` drops. + let handle = tokio::spawn(async move { + a.grow_to(2 * PERMIT_BYTES).await; + a + }); + tokio::time::sleep(Duration::from_millis(100)).await; + assert!(!handle.is_finished()); + handle + }; + drop(b); + let a = tokio::time::timeout(Duration::from_secs(5), grow) + .await + .unwrap() + .unwrap(); + assert_eq!(a.bytes(), 2 * PERMIT_BYTES); + assert_eq!(budget.reserved(), 2 * PERMIT_BYTES); + + let mut a = a; + a.grow_to(PERMIT_BYTES).await; // smaller: no-op + assert_eq!(a.bytes(), 2 * PERMIT_BYTES); + } +} diff --git a/crates/lance-context-core/src/metrics.rs b/crates/lance-context-core/src/metrics.rs index 2ee86eef..66bb67b2 100644 --- a/crates/lance-context-core/src/metrics.rs +++ b/crates/lance-context-core/src/metrics.rs @@ -65,6 +65,12 @@ pub const ROLLOUT_WAL_MERGE_ERRORS: &str = "rollout_wal_merge_errors_total"; /// before it was noticed in crash logs. pub const ROLLOUT_WAL_PENDING_GENERATIONS: &str = "rollout_wal_pending_generations"; +/// Time a merge spent waiting for its initial reservation from the +/// process-wide merge memory budget. Near zero when the budget is not the +/// bottleneck; a rising tail means merges are queueing on memory, which is the +/// intended behaviour under load rather than an error. +pub const ROLLOUT_MERGE_BUDGET_WAIT: &str = "rollout_merge_budget_wait_seconds"; + /// Emit a histogram sample in seconds. No-op without the `metrics` feature. #[cfg(feature = "metrics")] macro_rules! observe_duration { diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index ddc716f9..4f650c9a 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -77,6 +77,7 @@ use lance_index::mem_wal::ShardManifest; use serde_json::Value; use uuid::Uuid; +use crate::merge_budget::MergeMemoryBudget; use crate::rollout::RolloutRecord; use crate::store::{ column_as, column_as_optional, relationship_field, relationship_list_item_field, @@ -388,6 +389,10 @@ pub struct RolloutStoreOptions { /// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations` /// histogram is emitted regardless. pub pending_generations_warn: 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. + pub merge_budget: Option>, /// Shared Lance [`Session`] used to open this store's base dataset (and, /// transitively, every flushed MemWAL generation it reads — those inherit /// the base dataset's session). @@ -475,6 +480,7 @@ impl RolloutStore { merge_max_generations, merge_max_bytes, pending_generations_warn, + merge_budget, session, } = options; let mut base = StorageBase::open( @@ -486,6 +492,7 @@ impl RolloutStore { merge_max_generations, merge_max_bytes, pending_generations_warn, + merge_budget, session, schema: Arc::new(rollout_schema()), key_column: "id".to_string(), @@ -2503,6 +2510,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, session: None, schema: legacy_schema.clone(), key_column: "id".to_string(), @@ -2823,6 +2831,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -2868,6 +2877,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }; { @@ -2918,6 +2928,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -2970,6 +2981,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }; let instance_a = RolloutStore::open_with_options(&uri, options("rollout-0")) @@ -3192,6 +3204,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, ..Default::default() }, ) @@ -3255,6 +3268,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, ..Default::default() }, ) @@ -3356,6 +3370,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3415,6 +3430,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3455,6 +3471,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3592,6 +3609,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3679,6 +3697,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3727,6 +3746,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3797,6 +3817,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3835,6 +3856,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3894,6 +3916,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -3936,6 +3959,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -4436,6 +4460,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, ..Default::default() }, ) @@ -4535,6 +4560,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await @@ -4568,6 +4594,7 @@ mod tests { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, }, ) .await diff --git a/crates/lance-context-core/src/store.rs b/crates/lance-context-core/src/store.rs index 3dfc022e..5a0acb9a 100644 --- a/crates/lance-context-core/src/store.rs +++ b/crates/lance-context-core/src/store.rs @@ -33,6 +33,7 @@ use tokio::task::JoinHandle; use tracing::{error, info, warn}; use uuid::Uuid; +use crate::merge_budget::MergeMemoryBudget; use crate::record::{ ContextRecord, LifecycleQueryOptions, RecordFilters, RecordPatch, Relationship, RetrieveResult, SearchResult, StateMetadata, UpdateResult, UpsertResult, LIFECYCLE_ACTIVE, @@ -307,6 +308,10 @@ pub struct ContextStoreOptions { /// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations` /// histogram is emitted regardless. pub pending_generations_warn: 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. + pub merge_budget: Option>, /// Whether [`ContextStore::add`] seals the memtable before returning, so the /// rows it wrote are immediately readable. /// @@ -339,6 +344,7 @@ impl Default for ContextStoreOptions { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, // Read-your-write by default; see the field docs. seal_on_add: true, } @@ -621,6 +627,7 @@ impl ContextStore { merge_max_generations: options.merge_max_generations, merge_max_bytes: options.merge_max_bytes, pending_generations_warn: options.pending_generations_warn, + merge_budget: options.merge_budget.clone(), session: None, schema: Arc::new(arrow_schema.clone()), key_column: "id".to_string(), @@ -2245,6 +2252,7 @@ impl ContextStore { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, // A compactor never appends, so the seal mode is irrelevant to it; // deferring keeps it from ever emitting a generation. seal_on_add: false, diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index b0b70fce..5bdbb309 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -67,6 +67,7 @@ use lance_index::IndexType; use tracing::{info, warn}; use uuid::Uuid; +use crate::merge_budget::{MergeMemoryBudget, MergeReservation}; use crate::metrics::{ count, observe_duration, observe_phase, observe_value, timer_elapsed, timer_start, }; @@ -205,6 +206,10 @@ pub struct PreparedMerge { merged_paths: Vec, batches: Vec, merge_schema: Arc, + /// Bytes reserved from the process-wide merge budget for `batches`. + /// Released when this is dropped, i.e. right after the commit consumes + /// the batches. `None` when no budget is configured. + reservation: Option, } impl PreparedMerge { @@ -240,6 +245,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, + /// 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` + /// disables the bound. See [`crate::merge_budget`]. + pub merge_budget: Option>, /// Shared, capacity-bounded Lance session. `None` preserves Lance's /// per-open default (a fresh 6 GiB index + 1 GiB metadata session *per /// store*, which is the source of unbounded per-append RSS growth). @@ -303,6 +313,8 @@ pub(crate) struct StorageBase { merge_max_bytes: usize, /// Per-shard pending-generation count at which reads warn; `0` disables. pending_generations_warn: usize, + /// Process-wide merge byte budget; `None` means unbounded. + merge_budget: Option>, /// Timestamp of the last successful [`Self::compact`] on this handle. last_compaction: Option>, /// Number of successful compactions performed by this handle. @@ -360,6 +372,7 @@ impl StorageBase { merge_max_generations, merge_max_bytes, pending_generations_warn, + merge_budget, session, schema, key_column, @@ -391,6 +404,7 @@ impl StorageBase { merge_max_generations, merge_max_bytes, pending_generations_warn, + merge_budget, session, schema, key_column, @@ -416,6 +430,7 @@ impl StorageBase { merge_max_generations, merge_max_bytes, pending_generations_warn, + merge_budget, session, schema, key_column, @@ -443,6 +458,7 @@ 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), + merge_budget, last_compaction: None, total_compactions: 0, last_compaction_error: None, @@ -905,7 +921,7 @@ impl StorageBase { // The expensive phase: pull a budgeted prefix of generations out of object // storage. Buffered in memory, so this is the part that must not hold an // exclusive lock. - let (merged_generations, merged_paths, batches, merge_schema) = + let (merged_generations, merged_paths, batches, merge_schema, reservation) = observe_phase!("read", self.read_flushed_generations(manifest).await)?; Ok(Some(PreparedMerge { @@ -913,6 +929,7 @@ impl StorageBase { merged_paths, batches, merge_schema, + reservation, })) } @@ -953,6 +970,7 @@ impl StorageBase { merged_paths, batches, merge_schema, + reservation, } = prepared; // Several sweepers can prepare the same immutable generations under a @@ -978,6 +996,9 @@ impl StorageBase { )?; self.pinned_version = None; } + // The batches are consumed; give the bytes back before the manifest + // drain and directory deletes, which hold no merge data. + drop(reservation); // Reuse the shard's *current* epoch rather than claiming a new one: // claiming would fence our own live writer. `commit_update` still fails @@ -1059,7 +1080,35 @@ impl StorageBase { async fn read_flushed_generations( &self, manifest: &ShardManifest, - ) -> LanceResult<(HashSet, Vec, Vec, Arc)> { + ) -> LanceResult<( + HashSet, + Vec, + Vec, + Arc, + Option, + )> { + // Reserve the whole per-merge allowance from the process budget before + // reading anything. Taken in one acquire so a merge never holds part + // of what it needs while waiting for the rest (see `merge_budget`). + // With the count cap disabled and no byte cap, the allowance is the + // full budget: the merge may read everything pending, alone. + let mut reservation = match &self.merge_budget { + Some(budget) => { + let want = if self.merge_max_bytes == 0 { + budget.limit() + } else { + self.merge_max_bytes.min(budget.limit()) + }; + let wait = timer_start!(); + let reservation = budget.reserve(want).await; + observe_duration!( + crate::metrics::ROLLOUT_MERGE_BUDGET_WAIT, + timer_elapsed!(wait) + ); + Some(reservation) + } + None => None, + }; let base_uri = self.dataset.uri().trim_end_matches('/').to_string(); let mut merged_generations: HashSet = HashSet::new(); let mut merged_paths: Vec = Vec::new(); @@ -1099,6 +1148,13 @@ impl StorageBase { if batch.num_rows() > 0 { let batch = align_batch_to_schema(batch, merge_schema.clone())?; buffered_bytes = buffered_bytes.saturating_add(batch.get_array_memory_size()); + // Only an oversized first generation reads past the initial + // reservation (a generation is indivisible). Grow to cover + // it; the holder already has at least as much as anyone + // else could be waiting for, so this cannot deadlock. + if let Some(reservation) = reservation.as_mut() { + reservation.grow_to(buffered_bytes).await; + } current_batches.push(batch); } } @@ -1111,7 +1167,13 @@ impl StorageBase { } let batches = dedupe_merge_batches(generation_batches, &self.key_column, merge_schema.clone())?; - Ok((merged_generations, merged_paths, batches, merge_schema)) + Ok(( + merged_generations, + merged_paths, + batches, + merge_schema, + reservation, + )) } /// Merge prepared WAL rows into the base table by primary key. diff --git a/crates/lance-context-metrics/src/lib.rs b/crates/lance-context-metrics/src/lib.rs index bf0b2919..5b565344 100644 --- a/crates/lance-context-metrics/src/lib.rs +++ b/crates/lance-context-metrics/src/lib.rs @@ -60,6 +60,7 @@ const JOB_LATENCY_METRICS: &[&str] = &[ "rollout_wal_merge_duration_seconds", "rollout_wal_merge_request_duration_seconds", "rollout_wal_merge_lock_wait_seconds", + "rollout_merge_budget_wait_seconds", ]; /// Buckets (upper bounds, generation counts) for `rollout_wal_pending_generations`. diff --git a/crates/lance-context-server/src/config.rs b/crates/lance-context-server/src/config.rs index c4eb86b5..3185e589 100644 --- a/crates/lance-context-server/src/config.rs +++ b/crates/lance-context-server/src/config.rs @@ -62,6 +62,16 @@ pub struct ServerConfig { )] pub rollout_wal_pending_warn_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). + /// ROLLOUT_MERGE_MAX_BYTES bounds ONE merge; this bounds all of them + /// together, so the number of masters or concurrent tasks no longer + /// decides the worker's peak memory. A merge that cannot fit waits for + /// another to release instead of failing. Default 3 GiB; `0` disables. + #[arg(long, env = "ROLLOUT_MERGE_MEMORY_BYTES", default_value = "3221225472")] + pub rollout_merge_memory_bytes: 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 @@ -226,6 +236,7 @@ mod tests { .expect("the documented minimal invocation must parse"); assert_eq!(config.rollout_merge_max_bytes, 1024 * 1024 * 1024); assert_eq!(config.rollout_wal_pending_warn_generations, 256); + assert_eq!(config.rollout_merge_memory_bytes, 3 * 1024 * 1024 * 1024); assert_eq!(config.rollout_flush_interval_secs, 30); assert_eq!(config.rollout_cleanup_interval_secs, 0); } diff --git a/crates/lance-context-server/src/routes/datagen.rs b/crates/lance-context-server/src/routes/datagen.rs index 5895d519..6268a762 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), + merge_budget: state.merge_budget.clone(), cleanup_interval_secs: None, }; let store = DatagenStore::open_with_options(&uri, options) diff --git a/crates/lance-context-server/src/routes/generic.rs b/crates/lance-context-server/src/routes/generic.rs index 410b8452..bf487069 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), + merge_budget: state.merge_budget.clone(), session: None, seal_on_add: req.seal_on_add, }; diff --git a/crates/lance-context-server/src/routes/rollouts.rs b/crates/lance-context-server/src/routes/rollouts.rs index c9c9a160..de75dd10 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), + merge_budget: state.merge_budget.clone(), session: state.rollout_session.clone(), }; diff --git a/crates/lance-context-server/src/state.rs b/crates/lance-context-server/src/state.rs index 1697ed1d..5affabc7 100644 --- a/crates/lance-context-server/src/state.rs +++ b/crates/lance-context-server/src/state.rs @@ -7,8 +7,8 @@ use std::time::Duration; use lance_context_core::{ join_uri, validate_store_name, ContextStore, ContextStoreOptions, DatagenStore, - DatagenStoreOptions, GenericStore, GenericStoreOptions, RolloutRegistry, RolloutStore, - RolloutStoreOptions, Session, + DatagenStoreOptions, GenericStore, GenericStoreOptions, MergeMemoryBudget, RolloutRegistry, + RolloutStore, RolloutStoreOptions, Session, }; use lru::LruCache; use tokio::sync::{Mutex, OwnedMutexGuard, RwLock}; @@ -123,6 +123,9 @@ 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, + /// Process-wide merge memory budget shared by every store; `None` when + /// disabled. See `lance_context_core::merge_budget`. + pub merge_budget: 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, @@ -315,6 +318,8 @@ 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, + merge_budget: (config.rollout_merge_memory_bytes > 0) + .then(|| MergeMemoryBudget::new(config.rollout_merge_memory_bytes)), rollout_cleanup_interval_secs: config.rollout_cleanup_interval_secs, rollout_flush_interval_secs: config.rollout_flush_interval_secs, blob_budget, @@ -381,6 +386,7 @@ impl AppState { rollout_merge_max_generations: 8, rollout_merge_max_bytes: 1024 * 1024 * 1024, rollout_wal_pending_warn_generations: 256, + merge_budget: None, rollout_cleanup_interval_secs: 0, rollout_flush_interval_secs: 0, blob_budget: None, @@ -417,6 +423,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), + merge_budget: self.merge_budget.clone(), session: self.rollout_session.clone(), } } @@ -606,6 +613,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), + merge_budget: self.merge_budget.clone(), cleanup_interval_secs: None, } } @@ -749,6 +757,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), + 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 1668bbce..c49b199d 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, + merge_budget: None, cleanup_interval_secs: None, }; let store = LocalStore::open_with_options(uri, options) diff --git a/crates/lance-context/src/unified_generic.rs b/crates/lance-context/src/unified_generic.rs index 70ba02ae..c448618a 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, + merge_budget: None, session: None, seal_on_add, }; @@ -64,6 +65,7 @@ impl GenericStore { merge_max_generations: None, merge_max_bytes: None, pending_generations_warn: None, + merge_budget: None, session: None, seal_on_add, };