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
26 changes: 25 additions & 1 deletion crates/lance-context-core/src/generic_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<Option<(ShardManifestStore, ShardManifest, PreparedMerge)>> {
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<usize> {
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<usize> {
self.base.pending_wal_generations().await
Expand Down
27 changes: 20 additions & 7 deletions crates/lance-context-server/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
),
);
}
}))
}
Expand Down
29 changes: 27 additions & 2 deletions crates/lance-context-server/src/sweeper.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,13 +110,38 @@ impl Sweepable for Arc<RwLock<GenericStore>> {
}

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<usize, String> {
// 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())
}
}

Expand Down
Loading