diff --git a/crates/lance-context-core/src/generic_store.rs b/crates/lance-context-core/src/generic_store.rs index c835c5a6..873df0ec 100644 --- a/crates/lance-context-core/src/generic_store.rs +++ b/crates/lance-context-core/src/generic_store.rs @@ -25,13 +25,15 @@ use std::sync::Arc; use arrow_array::RecordBatch; use arrow_schema::{ArrowError, Schema}; use futures::TryStreamExt; +use lance::dataset::mem_wal::ShardManifestStore; use lance::dataset::optimize::CompactionMetrics; use lance::session::Session; 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::store::{CompactionConfig, CompactionStats}; -use crate::store_base::{ListSource, StorageBase, StorageBaseOptions}; +use crate::store_base::{ListSource, PreparedMerge, StorageBase, StorageBaseOptions}; use lance_context_api::schema_spec::{SchemaSpec, ID_COLUMN}; /// Schema-metadata key holding the serialized [`SchemaSpec`], so a store can be @@ -380,6 +382,28 @@ impl GenericStore { self.base.cleanup_own_shard().await } + /// The shared-lock half of [`Self::cleanup_wal`]: seal, then read a + /// budgeted prefix of flushed generations into memory. Callers holding a + /// read lock run this while appends continue, then take the write lock + /// only for [`Self::commit_prepared_merge`]. + pub async fn prepare_cleanup_merge( + &self, + ) -> LanceResult> { + self.base.prepare_cleanup_merge().await + } + + /// Commit a merge prepared by [`Self::prepare_cleanup_merge`]. + pub async fn commit_prepared_merge( + &mut self, + manifest_store: &ShardManifestStore, + manifest: &ShardManifest, + prepared: PreparedMerge, + ) -> LanceResult { + self.base + .commit_prepared_merge(manifest_store, manifest, prepared) + .await + } + /// Generations pending merge across all shards. Read-only. pub async fn pending_wal_generations(&self) -> LanceResult { self.base.pending_wal_generations().await diff --git a/crates/lance-context-server/src/state.rs b/crates/lance-context-server/src/state.rs index 9857e0d3..ef6d15fb 100644 --- a/crates/lance-context-server/src/state.rs +++ b/crates/lance-context-server/src/state.rs @@ -908,13 +908,26 @@ impl AppState { return; }; // Every kind, not just rollout: generations accumulate in each - // one, and only the base table absorbs them. - sweeper::merge_pass(sweeper::resident(&state.rollout_stores).await, pass_timeout) - .await; - sweeper::merge_pass(sweeper::resident(&state.datagen_stores).await, pass_timeout) - .await; - sweeper::merge_pass(sweeper::resident(&state.generic_stores).await, pass_timeout) - .await; + // one, and only the base table absorbs them. The kinds run + // concurrently: each pass is a serial walk with a per-store + // timeout of several minutes, so with hundreds of resident + // rollout stores a serial kind order let generic stores go + // unswept for hours while their generations piled into the + // tens of thousands and every read re-opened all of them. + tokio::join!( + sweeper::merge_pass( + sweeper::resident(&state.rollout_stores).await, + pass_timeout + ), + sweeper::merge_pass( + sweeper::resident(&state.datagen_stores).await, + pass_timeout + ), + sweeper::merge_pass( + sweeper::resident(&state.generic_stores).await, + pass_timeout + ), + ); } })) } diff --git a/crates/lance-context-server/src/sweeper.rs b/crates/lance-context-server/src/sweeper.rs index 0500e257..b629cfa0 100644 --- a/crates/lance-context-server/src/sweeper.rs +++ b/crates/lance-context-server/src/sweeper.rs @@ -110,13 +110,38 @@ impl Sweepable for Arc> { } async fn flush(&self) -> Result<(), String> { + // Same shape as rollout: generic stores default to a deferred seal and + // take the same one-row-per-append traffic, so the count-triggered + // merge rides this timer too. Without it a hot generic store depends + // entirely on the slower cleanup sweeper reaching it. let guard = self.read().await; - guard.flush().await.map_err(|e| e.to_string()) + let result = guard.flush().await.map_err(|e| e.to_string()); + if result.is_ok() { + drop(guard); + let mut guard = self.write().await; + guard.maybe_merge_wal().await.map_err(|e| e.to_string())?; + } + result } async fn merge_wal(&self) -> Result { + // Prepare under the shared lock so appends and reads keep running while + // generations are read; hold the exclusive lock only for the commit. + let prepared = { + let guard = self.read().await; + guard + .prepare_cleanup_merge() + .await + .map_err(|e| e.to_string())? + }; + let Some((manifest_store, manifest, prepared)) = prepared else { + return Ok(0); + }; let mut guard = self.write().await; - guard.cleanup_wal().await.map_err(|e| e.to_string()) + guard + .commit_prepared_merge(&manifest_store, &manifest, prepared) + .await + .map_err(|e| e.to_string()) } }