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
73 changes: 42 additions & 31 deletions crates/lance-context-master/src/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,10 @@ use lance_context_api::ExperimentSummary;
/// scan round.
const OBSERVE_TIMEOUT: Duration = Duration::from_secs(30);

/// Bound on one steady-state `_stats` maintenance pass so a slow object store
/// cannot wedge the scanner loop.
/// Bound on the compaction half of one `_stats` maintenance pass, which runs
/// under the store lock, so a slow object store cannot wedge the scanner loop.
///
/// Deliberately not applied to the first pass after startup: see
/// [`maintain_stats`].
/// The cleanup half is not bounded: see [`maintain_stats`].
const MAINTENANCE_TIMEOUT: Duration = Duration::from_secs(300);

/// Bound on the merge and compaction steps of retiring one experiment. Larger
Expand All @@ -47,58 +46,70 @@ const RETIRE_TIMEOUT: Duration = Duration::from_secs(600);
/// one per round and this pass is cheap. It remains necessary to reclaim those
/// versions, and to recover deployments that ran the old per-row path.
///
/// # The first pass runs without a timeout
/// # The store lock is held only for the compaction
///
/// A deployment upgraded from the per-row path can arrive with a chain
/// hundreds of thousands of versions long (246k+ observed). Compacting and
/// pruning that cannot finish inside `MAINTENANCE_TIMEOUT`, so every pass timed
/// out, rolled back, and left the table exactly as bloated as before — the
/// bound guaranteed the table could never recover. The first pass after startup
/// therefore runs unbounded, and subsequent passes take the bound: by then the
/// backlog is gone and any pass exceeding it is a genuine fault.
/// Compaction rewrites the table and must exclude other writers, so it runs
/// under `state.stats`. The cleanup that follows deletes objects no live
/// version references; it needs no exclusive access, only time -- a backlog of
/// tens of thousands of versions is tens of thousands of deletes against a
/// remote object store. It used to run under the same lock, and every sweep
/// and every post-compaction stats refresh on every master waited behind it
/// (36 minutes observed after a routine restart). It now runs on a detached
/// handle with the lock released; the store is re-opened afterwards so its
/// in-memory handle never points at a pruned manifest.
///
/// The cleanup is deliberately unbounded. Lance deletes as it goes, so a
/// timeout would only abandon the pass mid-way, not undo it, and with the lock
/// released nothing waits on it any more. Only the compaction is bounded by
/// [`MAINTENANCE_TIMEOUT`].
///
/// Callers must hold the `stats-writer` coordination lock so only one replica
/// ever rewrites the dataset.
pub async fn maintain_stats(state: &Arc<MasterState>) -> lance::Result<()> {
let ttl = Duration::from_secs(state.config.stats_history_ttl_secs);
let start = std::time::Instant::now();
let mut stats = state.stats.lock().await;

// `swap` so exactly one pass per process is unbounded, even if several
// scanner ticks race here.
let first_pass = state.stats_maintenance_done.swap(true, Ordering::SeqCst);
let outcome = if first_pass {
tokio::time::timeout(MAINTENANCE_TIMEOUT, stats.maintain(ttl))
let compacted = {
let mut stats = state.stats.lock().await;
tokio::time::timeout(MAINTENANCE_TIMEOUT, stats.compact())
.await
.unwrap_or_else(|_| Err(lance::Error::io("stats maintenance timed out")))
} else {
tracing::info!(
version = stats.version(),
"running first stats maintenance pass without a timeout; \
a table carried over from the per-row write path can take a while to reclaim"
);
stats.maintain(ttl).await
.unwrap_or_else(|_| Err(lance::Error::io("stats compaction timed out")))
};
let outcome = match compacted {
Ok((compaction, cleaner)) => {
let removal = cleaner.cleanup(ttl).await;
// Re-open regardless of the cleanup's outcome: a failed cleanup
// may still have deleted manifests the current handle references.
let reload = state.stats.lock().await.reload().await;
match (removal, reload) {
(Ok(removal), Ok(())) => Ok((compaction, removal)),
(Err(e), _) | (Ok(_), Err(e)) => Err(e),
}
}
Err(e) => Err(e),
};
let version = state.stats.lock().await.version();

match outcome {
Ok((compaction, removal)) => {
state.stats_maintenance_failures.store(0, Ordering::Relaxed);
state
.stats_last_reclaimed_version
.store(stats.version(), Ordering::Relaxed);
.store(version, Ordering::Relaxed);
metrics::histogram!("master_stats_maintenance_duration_seconds")
.record(start.elapsed().as_secs_f64());
metrics::counter!("master_stats_versions_removed_total")
.increment(removal.old_versions);
metrics::gauge!("master_stats_version").set(stats.version() as f64);
metrics::gauge!("master_stats_version").set(version as f64);
metrics::gauge!("master_stats_maintenance_consecutive_failures").set(0.0);
metrics::gauge!("master_stats_unreclaimed_versions").set(0.0);
tracing::info!(
fragments_removed = compaction.fragments_removed,
fragments_added = compaction.fragments_added,
old_versions_removed = removal.old_versions,
bytes_removed = removal.bytes_removed,
version = stats.version(),
version,
elapsed_secs = start.elapsed().as_secs(),
"stats maintenance complete"
);
Ok(())
Expand All @@ -121,7 +132,7 @@ pub async fn maintain_stats(state: &Arc<MasterState>) -> lance::Result<()> {
.fetch_add(1, Ordering::Relaxed)
+ 1;
let last_reclaimed = state.stats_last_reclaimed_version.load(Ordering::Relaxed);
let unreclaimed = stats.version().saturating_sub(last_reclaimed);
let unreclaimed = version.saturating_sub(last_reclaimed);

metrics::counter!("master_stats_maintenance_failures_total").increment(1);
metrics::gauge!("master_stats_maintenance_consecutive_failures").set(failures as f64);
Expand All @@ -131,7 +142,7 @@ pub async fn maintain_stats(state: &Arc<MasterState>) -> lance::Result<()> {
error = %e,
consecutive_failures = failures,
unreclaimed_versions = unreclaimed,
version = stats.version(),
version,
"stats maintenance failed; old manifests are not being reclaimed"
);
Err(e)
Expand Down
27 changes: 24 additions & 3 deletions crates/lance-context-master/src/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@ use crate::task_store::{TaskClaim, TaskKinds};
/// anything still over the threshold is picked up by the next tick.
const MAX_SWEEP_ENQUEUE: usize = 256;

/// How long a finished compaction waits for the `stats-writer` lock to refresh
/// its stats row before giving up and leaving it to the next scan round.
const STATS_REFRESH_LOCK_WAIT: Duration = Duration::from_secs(10);

/// Task-target prefix marking a generic store. Store names match
/// `[A-Za-z0-9_][A-Za-z0-9._-]*`, so a `:` can never appear in a bare name
/// and the prefix is unambiguous. Carrying the kind in the target string --
Expand Down Expand Up @@ -383,13 +387,30 @@ async fn merge_wal_one(http: &reqwest::Client, url: &str) -> Result<WorkerMerge,

/// Refresh the stats row for `name` after a successful compaction: re-observe
/// fragment/row counts and bump `last_compaction`/`total_compactions`.
///
/// Best-effort. The compaction itself is already committed, and the next scan
/// round refreshes the row anyway, so this never waits long for the
/// `stats-writer` lock: a holder mid-maintenance would otherwise pin this
/// task's concurrency slot on every master for as long as it runs.
async fn update_stats_after_compaction(state: &Arc<MasterState>, name: &str, store: &RolloutStore) {
let guard = match state.task_store.coordination_lock("stats-writer").await {
Ok(guard) => guard,
Err(e) => {
let guard = match tokio::time::timeout(
STATS_REFRESH_LOCK_WAIT,
state.task_store.coordination_lock("stats-writer"),
)
.await
{
Ok(Ok(guard)) => guard,
Ok(Err(e)) => {
tracing::warn!(store = %name, error = %e, "stats writer lock failed");
return;
}
Err(_) => {
tracing::info!(
store = %name,
"stats writer busy; leaving the post-compaction refresh to the next scan"
);
return;
}
};
let obs = match store.observe().await {
Ok(obs) => obs,
Expand Down
8 changes: 0 additions & 8 deletions crates/lance-context-master/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,13 +104,6 @@ pub struct MasterState {
pub http: reqwest::Client,
/// Process-wide compaction permits shared by scheduler and retirement work.
pub(crate) compaction_permits: Arc<Semaphore>,
/// Whether this process has already run one `_stats` maintenance pass.
///
/// The first pass runs without a timeout so a deployment carrying a version
/// chain from the old per-row write path can actually reclaim it; a bounded
/// first pass times out forever and never recovers. See
/// [`crate::scanner::maintain_stats`].
pub stats_maintenance_done: std::sync::atomic::AtomicBool,
/// Consecutive `_stats` maintenance failures, for alerting.
///
/// Failure was previously silent: the success counter simply stopped
Expand Down Expand Up @@ -167,7 +160,6 @@ impl MasterState {
task_store,
http: reqwest::Client::new(),
compaction_permits: Arc::new(Semaphore::new(compaction_concurrency)),
stats_maintenance_done: std::sync::atomic::AtomicBool::new(false),
stats_maintenance_failures: std::sync::atomic::AtomicU64::new(0),
stats_last_reclaimed_version: std::sync::atomic::AtomicU64::new(0),
});
Expand Down
85 changes: 79 additions & 6 deletions crates/lance-context-master/src/stats_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -544,10 +544,23 @@ impl StatsStore {
///
/// Callers must serialize this with other mutations (it takes `&mut self`)
/// and, across master replicas, hold the `stats-writer` coordination lock.
/// Callers that share this store behind a lock should prefer
/// [`Self::compact`] followed by [`Self::cleanup`] on the returned handle,
/// so the lock is not held for the cleanup's many object-store deletes.
pub async fn maintain(
&mut self,
older_than: Duration,
) -> LanceResult<(CompactionMetrics, RemovalStats)> {
let (compaction, cleaner) = self.compact().await?;
let removal = cleaner.cleanup(older_than).await?;
self.reload().await?;
Ok((compaction, removal))
}

/// The compaction half of [`Self::maintain`]: rewrite the table's fragments
/// into one. Returns a [`StatsCleaner`] that prunes old versions *without*
/// borrowing the store, so a caller can release its lock first.
pub async fn compact(&mut self) -> LanceResult<(CompactionMetrics, StatsCleaner)> {
self.dataset.checkout_latest().await?;

let options = CompactionOptions {
Expand All @@ -560,16 +573,45 @@ impl StatsStore {
};
let compaction = compact_files(&mut self.dataset, options, None).await?;

// Re-open so the handle (and the cleanup below) sees the rewritten
// version rather than the pre-compaction manifest.
// Re-open so the handle (and the cleanup) sees the rewritten version
// rather than the pre-compaction manifest.
self.reload().await?;
Ok((
compaction,
StatsCleaner {
dataset: self.dataset.clone(),
},
))
}

/// Re-open the dataset at its latest version. Call after a
/// [`StatsCleaner::cleanup`] so the in-memory handle does not reference
/// manifests the cleanup removed.
pub async fn reload(&mut self) -> LanceResult<()> {
self.dataset = Self::load(&self.uri, self.storage_options.clone()).await?;
Ok(())
}
}

/// A detached handle for the old-version cleanup half of maintenance.
///
/// Cleanup only reads manifests and deletes objects no live version references,
/// so it needs no exclusive access to the store: readers on the current
/// version, and even writers committing new versions, are unaffected. What it
/// does need is time -- tens of thousands of deletes against a remote object
/// store -- and holding the store's lock for that stalled every sweep and
/// every post-compaction stats refresh on every master for the duration.
pub struct StatsCleaner {
dataset: Dataset,
}

impl StatsCleaner {
/// Drop manifest versions older than `older_than` and the files only they
/// referenced.
pub async fn cleanup(self, older_than: Duration) -> LanceResult<RemovalStats> {
let grace = chrono::TimeDelta::from_std(older_than)
.map_err(|e| LanceError::io(format!("invalid stats history TTL: {e}")))?;
let removal = self.dataset.cleanup_old_versions(grace, None, None).await?;
self.dataset = Self::load(&self.uri, self.storage_options.clone()).await?;

Ok((compaction, removal))
self.dataset.cleanup_old_versions(grace, None, None).await
}
}

Expand Down Expand Up @@ -695,6 +737,37 @@ mod tests {
assert_eq!(s.list(None, 10, 0).await.unwrap().len(), 1);
}

/// The cleanup half of maintenance runs on a detached handle so the master
/// can release its store lock first. The store must stay fully usable --
/// reads and writes -- while that cleanup runs, and be intact after it.
#[tokio::test]
async fn cleanup_runs_detached_from_the_store() {
let dir = TempDir::new().unwrap();
let mut s = new_store(&dir).await;
for i in 0..20 {
s.upsert(&sample("exp-a", i)).await.unwrap();
s.upsert(&sample("exp-b", i)).await.unwrap();
}
let (compaction, cleaner) = s.compact().await.unwrap();
assert!(compaction.fragments_removed > 0, "nothing compacted");

// The store is free while the cleaner exists: write through it.
s.upsert(&sample("exp-c", 1)).await.unwrap();
assert_eq!(s.count(None).await.unwrap(), 3);

let removal = cleaner.cleanup(Duration::from_secs(0)).await.unwrap();
assert!(removal.old_versions > 0, "no versions reclaimed");

// The write that landed during cleanup is the live version; reloading
// onto it keeps every row.
s.reload().await.unwrap();
let rows = s.list(None, 10, 0).await.unwrap();
assert_eq!(rows.len(), 3);
assert_eq!(rows[0].row_count, 19);
s.upsert(&sample("exp-d", 1)).await.unwrap();
assert_eq!(s.count(None).await.unwrap(), 4);
}

#[tokio::test]
async fn no_compaction_sentinel_maps_to_none() {
let row = sample("x", 1);
Expand Down
Loading