Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions crates/lance-context-core/src/datagen_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand Down Expand Up @@ -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<usize>,
/// 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<Arc<MergeMemoryBudget>>,
/// Periodically merge this writer's pending generations. `None` or zero
/// disables the timer.
pub cleanup_interval_secs: Option<u64>,
Expand Down Expand Up @@ -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
Expand Down
87 changes: 87 additions & 0 deletions crates/lance-context-core/src/generic_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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<usize>,
/// 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<Arc<MergeMemoryBudget>>,
/// Shared, capacity-bounded Lance session.
pub session: Option<Arc<Session>>,
/// Whether [`GenericStore::add`] seals before returning, making the rows it
Expand Down Expand Up @@ -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`
Expand Down Expand Up @@ -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();
Expand Down
2 changes: 2 additions & 0 deletions crates/lance-context-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
239 changes: 239 additions & 0 deletions crates/lance-context-core/src/merge_budget.rs
Original file line number Diff line number Diff line change
@@ -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<MergeMemoryBudget>,
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<Self> {
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<Self>, 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);
}
}
Loading
Loading