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
154 changes: 154 additions & 0 deletions crates/lance-context-core/src/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,14 @@

use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::time::Duration;

use arrow_array::{Int64Array, RecordBatch, RecordBatchIterator, StringArray};
use arrow_schema::{ArrowError, DataType, Field, Schema};
use chrono::Utc;
use futures::TryStreamExt;
use lance::dataset::cleanup::RemovalStats;
use lance::dataset::optimize::{compact_files, CompactionMetrics, CompactionOptions};
use lance::dataset::{builder::DatasetBuilder, Dataset, WriteMode, WriteParams};
use lance::io::{ObjectStoreParams, StorageOptionsAccessor};
use lance::{Error as LanceError, Result as LanceResult};
Expand Down Expand Up @@ -347,6 +350,73 @@ impl RolloutRegistry {
pub fn uri(&self) -> &str {
&self.uri
}

/// Current dataset version (manifest chain head).
pub fn version(&self) -> u64 {
self.dataset.version().version
}

/// Fold the one-row append fragments produced by [`Self::upsert`] into a
/// few, then drop manifest versions older than `older_than`.
///
/// Every `create` is a delete plus an append, so the registry gains two
/// versions and one fragment per store. Lance keeps every manifest until
/// cleaned, and every manifest lists every fragment, so an unmaintained
/// registry grows quadratically: 34k versions of ~840 KB each (29 GB of
/// manifests) for a 10k-row, three-column table was observed in
/// production, and each `contains`/`get` re-reads the head manifest.
///
/// Callers that share this registry behind a lock should prefer
/// [`Self::compact`] followed by [`RegistryCleaner::cleanup`] so the lock
/// is not held across the cleanup's 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`]. Returns a
/// [`RegistryCleaner`] that prunes old versions without borrowing the
/// registry.
pub async fn compact(&mut self) -> LanceResult<(CompactionMetrics, RegistryCleaner)> {
self.reload().await?;
let options = CompactionOptions {
target_rows_per_fragment: 1_048_576,
materialize_deletions: true,
materialize_deletions_threshold: 0.0,
..Default::default()
};
let compaction = compact_files(&mut self.dataset, options, None).await?;
self.reload().await?;
Ok((
compaction,
RegistryCleaner {
dataset: self.dataset.clone(),
},
))
}
}

/// A detached handle for the old-version cleanup half of registry
/// maintenance. See `StatsCleaner` in the master for why this is split off:
/// cleanup deletes objects no live version references, so it needs no
/// exclusive access, only time.
pub struct RegistryCleaner {
dataset: Dataset,
}

impl RegistryCleaner {
/// 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 registry history TTL: {e}")))?;
self.dataset.cleanup_old_versions(grace, None, None).await
}
}

#[cfg(test)]
Expand All @@ -361,6 +431,90 @@ mod tests {
.unwrap()
}

/// Maintenance folds the per-upsert fragments and prunes old manifests
/// while keeping every live row; the registry keeps working afterwards.
#[tokio::test]
async fn maintain_bounds_versions_and_preserves_rows() {
let dir = TempDir::new().unwrap();
let mut r = new_registry(&dir).await;
for i in 0..20 {
r.upsert(&format!("exp-{i}"), &format!("/data/exp-{i}.lance"))
.await
.unwrap();
}
let before = r.version();
assert!(
before >= 40,
"expected two versions per upsert, got {before}"
);

let (compaction, removal) = r.maintain(Duration::from_secs(0)).await.unwrap();
assert!(compaction.fragments_removed > 0, "nothing compacted");
assert!(removal.old_versions > 0, "no versions reclaimed");

let mut names: Vec<String> = r
.list()
.await
.unwrap()
.into_iter()
.map(|e| e.name)
.collect();
names.sort();
assert_eq!(names.len(), 20);
assert_eq!(names[0], "exp-0");
assert!(r.contains("exp-19").await.unwrap());
r.upsert("exp-new", "/data/exp-new.lance").await.unwrap();
assert!(r.contains("exp-new").await.unwrap());
}

/// A second handle on the same URI (a worker) keeps creating stores while
/// another (the master) maintains. Nothing the worker wrote is lost and
/// both handles see every row afterwards.
#[tokio::test]
async fn maintenance_does_not_lose_concurrent_worker_upserts() {
let dir = TempDir::new().unwrap();
let uri = dir.path().join("_registry.lance");
let uri = uri.to_str().unwrap();
let mut master = RolloutRegistry::open_or_create(uri, None).await.unwrap();
let mut worker = RolloutRegistry::open_or_create(uri, None).await.unwrap();
for i in 0..20 {
worker.upsert(&format!("old-{i}"), "/x").await.unwrap();
}

let (_, cleaner) = master.compact().await.unwrap();
// Worker writes land between the compaction and the cleanup ...
for i in 0..5 {
worker.upsert(&format!("mid-{i}"), "/x").await.unwrap();
}
cleaner.cleanup(Duration::from_secs(0)).await.unwrap();
master.reload().await.unwrap();
// ... and after it.
worker.upsert("late", "/x").await.unwrap();

assert_eq!(master.list().await.unwrap().len(), 26);
assert_eq!(worker.list().await.unwrap().len(), 26);
assert!(master.contains("mid-3").await.unwrap());
assert!(worker.contains("late").await.unwrap());
}

/// The cleanup half runs on a detached handle: the registry stays usable
/// (reads and writes) while the cleaner is outstanding.
#[tokio::test]
async fn cleanup_runs_detached_from_the_registry() {
let dir = TempDir::new().unwrap();
let mut r = new_registry(&dir).await;
for i in 0..20 {
r.upsert(&format!("exp-{i}"), "/x").await.unwrap();
}
let (_, cleaner) = r.compact().await.unwrap();
r.upsert("during", "/x").await.unwrap();
assert!(r.contains("during").await.unwrap());
let removal = cleaner.cleanup(Duration::from_secs(0)).await.unwrap();
assert!(removal.old_versions > 0);
r.reload().await.unwrap();
assert_eq!(r.list().await.unwrap().len(), 21);
}

#[tokio::test]
async fn create_or_load_recovers_when_another_caller_wins() {
let dir = TempDir::new().unwrap();
Expand Down
60 changes: 59 additions & 1 deletion crates/lance-context-master/src/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,8 @@ use std::time::Duration;
use chrono::Utc;
use futures::stream::{self, StreamExt};
use lance_context_core::{
CompactionConfig, GenericStore, GenericStoreOptions, RolloutStore, RolloutStoreOptions,
CompactionConfig, GenericStore, GenericStoreOptions, RolloutRegistry, RolloutStore,
RolloutStoreOptions,
};
use tokio::sync::Semaphore;
use tokio::task::JoinHandle;
Expand Down Expand Up @@ -150,6 +151,55 @@ pub async fn maintain_stats(state: &Arc<MasterState>) -> lance::Result<()> {
}
}

/// Compact a store registry and prune its old manifest versions.
///
/// Workers write the registry on every store create/delete as a delete plus an
/// append, so it gains two versions and one fragment per store and, without
/// this, every manifest lists every fragment: 34k versions of ~840 KB each
/// were observed for a 10k-row table, and every `create`/`get` re-read one.
///
/// Same shape as [`maintain_stats`]: the compaction (a `Rewrite` that must not
/// race another master's) runs under `state.stats`-style exclusion via the
/// registry lock, bounded by [`MAINTENANCE_TIMEOUT`]; the cleanup runs on a
/// detached handle with the lock released. Workers keep appending in between:
/// Lance retries their commits past the rewrite, and the grace window keeps
/// any version a worker may still be reading.
async fn maintain_registry(
state: &Arc<MasterState>,
label: &'static str,
registry: &tokio::sync::RwLock<RolloutRegistry>,
) -> lance::Result<()> {
let ttl = Duration::from_secs(state.config.stats_history_ttl_secs);
let start = std::time::Instant::now();
let (compaction, cleaner) = {
let mut registry = registry.write().await;
tokio::time::timeout(MAINTENANCE_TIMEOUT, registry.compact())
.await
.unwrap_or_else(|_| Err(lance::Error::io("registry compaction timed out")))?
};
let removal = cleaner.cleanup(ttl).await;
let reload = registry.write().await.reload().await;
let removal = match (removal, reload) {
(Ok(removal), Ok(())) => removal,
(Err(e), _) | (Ok(_), Err(e)) => return Err(e),
};
let version = registry.read().await.version();
metrics::counter!("master_registry_versions_removed_total", "registry" => label)
.increment(removal.old_versions);
metrics::gauge!("master_registry_version", "registry" => label).set(version as f64);
tracing::info!(
registry = label,
fragments_removed = compaction.fragments_removed,
fragments_added = compaction.fragments_added,
old_versions_removed = removal.old_versions,
bytes_removed = removal.bytes_removed,
version,
elapsed_secs = start.elapsed().as_secs(),
"registry maintenance complete"
);
Ok(())
}

/// Run a single scan pass: refresh every experiment's stats row and drop rows
/// for experiments no longer in the registry. Returns the number of
/// experiments successfully observed.
Expand Down Expand Up @@ -179,6 +229,14 @@ async fn try_scan_once(state: &Arc<MasterState>, maintain: bool) -> lance::Resul
if let Err(e) = maintain_stats(state).await {
tracing::warn!(error = %e, "stats maintenance failed");
}
for (label, registry) in [
("rollout", &state.registry),
("generic", &state.generic_registry),
] {
if let Err(e) = maintain_registry(state, label, registry).await {
tracing::warn!(registry = label, error = %e, "registry maintenance failed");
}
}
}
let release = state.task_store.release_coordination_lock(guard).await;
match (result, release) {
Expand Down
Loading