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
7 changes: 7 additions & 0 deletions crates/lance-context-master/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,13 @@ pub struct MasterConfig {
#[arg(long, env = "COMPACTION_MAX_SOURCE_FRAGMENTS", default_value_t = 32)]
pub compaction_max_source_fragments: usize,

/// Rebuild the `id` ZoneMap index after every compaction that rewrote
/// fragments, as a dependent `IndexId` task. Compaction replaces the
/// fragments the index describes, so without this the index only ever
/// covers whatever was on disk when it was last built by hand.
#[arg(long, env = "INDEX_AFTER_COMPACTION", default_value_t = true, action = clap::ArgAction::Set)]
pub index_after_compaction: bool,

/// Maximum bytes per compacted output file. `0` uses Lance's default.
#[arg(
long,
Expand Down
1 change: 1 addition & 0 deletions crates/lance-context-master/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -679,6 +679,7 @@ mod tests {
compaction_threads: 1,
compaction_batch_size: 8,
compaction_max_source_fragments: 32,
index_after_compaction: false,
compaction_max_bytes_per_file: 1024 * 1024 * 1024,
merge_wal_interval_secs: 0,
merge_wal_min_generations: 8,
Expand Down
111 changes: 102 additions & 9 deletions crates/lance-context-master/src/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -143,7 +143,7 @@ async fn run_task(state: &Arc<MasterState>, claim: TaskClaim, timing: TaskClaimT

let started = std::time::Instant::now();
let outcome = match task.kind {
TaskKind::Compact => run_compaction(state, &task.target).await,
TaskKind::Compact => run_compaction(state, &task).await,
TaskKind::MergeWal => run_merge_wal(state, &task.target).await,
TaskKind::IndexId => run_index_id(state, &task.target).await,
};
Expand Down Expand Up @@ -174,15 +174,39 @@ async fn run_task(state: &Arc<MasterState>, claim: TaskClaim, timing: TaskClaimT

/// Compact one experiment. The task-store claim owns the per-experiment write
/// lock for the full execution.
async fn run_compaction(state: &Arc<MasterState>, target: &str) -> Result<String, String> {
let (kind, name) = parse_target(target);
///
/// A compaction that rewrote fragments invalidates the `id` ZoneMap index
/// (per-fragment min/max), so it enqueues an [`TaskKind::IndexId`] that
/// depends on this task. The dependency keeps the two from contending for the
/// per-target lock and lets the index wait its turn behind the queue. Skipped
/// when nothing was rewritten and when `index_after_compaction` is off.
async fn run_compaction(state: &Arc<MasterState>, task: &TaskRecord) -> Result<String, String> {
let (kind, name) = parse_target(&task.target);
if kind != StoreKind::Rollout {
// Base-table compaction of generic stores is not scheduled by the
// master yet; only WAL merges are. Refuse rather than open the store
// through the rollout code path with the wrong URI and schema.
return Err(format!("compaction is not scheduled for {kind:?} stores"));
}
compact_inner(state, name).await
let metrics = compact_inner(state, name).await?;
if state.config.index_after_compaction && metrics.fragments_added > 0 {
// Best-effort: the compaction itself is done; the next compaction of
// this target re-enqueues the index anyway.
if let Err(error) = enqueue_with_deps(
state,
TaskKind::IndexId,
&task.target,
vec![task.id.clone()],
)
.await
{
tracing::warn!(target = %task.target, %error, "failed to enqueue post-compaction id index");
}
}
Ok(format!(
"removed {} / added {} fragments",
metrics.fragments_removed, metrics.fragments_added
))
}

/// Build a ZoneMap scalar index on one experiment's `id` column. Shares the
Expand All @@ -209,7 +233,10 @@ async fn index_id_inner(state: &Arc<MasterState>, name: &str) -> Result<String,
Ok("built zonemap index on id".to_string())
}

async fn compact_inner(state: &Arc<MasterState>, name: &str) -> Result<String, String> {
async fn compact_inner(
state: &Arc<MasterState>,
name: &str,
) -> Result<lance::dataset::optimize::CompactionMetrics, String> {
let uri = state.rollout_uri(name);
let config = state.compaction_config();
let _permit = state
Expand All @@ -227,10 +254,7 @@ async fn compact_inner(state: &Arc<MasterState>, name: &str) -> Result<String, S
.await
.map_err(|e| e.to_string())?;
update_stats_after_compaction(state, name, &store).await;
Ok(format!(
"removed {} / added {} fragments",
metrics.fragments_removed, metrics.fragments_added
))
Ok(metrics)
}

/// Shape of the worker's merge-wal response (`{ "reclaimed": n }`).
Expand Down Expand Up @@ -737,6 +761,7 @@ mod tests {
compaction_threads: 1,
compaction_batch_size: 8,
compaction_max_source_fragments: 32,
index_after_compaction: false,
compaction_max_bytes_per_file: 1024 * 1024 * 1024,
merge_wal_interval_secs: 0,
merge_wal_min_generations: 2,
Expand Down Expand Up @@ -828,6 +853,74 @@ mod tests {
worker.abort();
}

/// A compaction that rewrote fragments enqueues an `IndexId` that depends
/// on it, and that index runs to Done after the compaction.
#[tokio::test]
#[ignore = "requires ETCD_TEST_ENDPOINTS"]
async fn compaction_that_rewrites_fragments_enqueues_id_index() {
let dir = TempDir::new().unwrap();
let mut cfg = config(&dir);
cfg.index_after_compaction = true;
let state = MasterState::new(cfg).await.unwrap();
let worker = spawn_scheduler(&state);

let name = "exp";
let uri = state.rollout_uri(name);
{
let mut store = RolloutStore::open(&uri).await.unwrap();
for i in 0..4 {
let rec = rollout_record(&format!("r{i}"));
store.add(&[rec]).await.unwrap();
store.cleanup_own_shard().await.unwrap();
}
}
state
.registry
.write()
.await
.upsert(name, &uri)
.await
.unwrap();
crate::scanner::scan_once(&state).await.unwrap();

let compact = enqueue(&state, TaskKind::Compact, name).await.unwrap();
assert_eq!(
await_terminal(&state, &compact.id).await.state,
TaskState::Done
);

let index = state
.task_store
.list()
.await
.unwrap()
.into_iter()
.find(|task| task.kind == TaskKind::IndexId && task.target == name)
.expect("compaction enqueued an IndexId task");
assert_eq!(index.depends_on, vec![compact.id.clone()]);
let status = await_terminal(&state, &index.id).await;
assert_eq!(status.state, TaskState::Done, "got {status:?}");
assert_eq!(status.detail.as_deref(), Some("built zonemap index on id"));

// A second compaction with nothing to rewrite does not enqueue another.
let again = enqueue(&state, TaskKind::Compact, name).await.unwrap();
assert_eq!(
await_terminal(&state, &again.id).await.state,
TaskState::Done
);
let indexes = state
.task_store
.list()
.await
.unwrap()
.into_iter()
.filter(|task| task.kind == TaskKind::IndexId)
.count();
assert_eq!(indexes, 1);

worker.abort();
}

/// Manual enqueue of an `IndexId` task -> dispatcher builds the ZoneMap
/// index -> task reaches Done with the expected detail summary.
#[tokio::test]
Expand Down
1 change: 1 addition & 0 deletions crates/lance-context-master/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -252,6 +252,7 @@ mod tests {
compaction_threads: 1,
compaction_batch_size: 8,
compaction_max_source_fragments: 32,
index_after_compaction: false,
compaction_max_bytes_per_file: 1024 * 1024 * 1024,
merge_wal_interval_secs: 0,
merge_wal_min_generations: 8,
Expand Down
1 change: 1 addition & 0 deletions crates/lance-context-master/src/task_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1251,6 +1251,7 @@ mod tests {
compaction_threads: 1,
compaction_batch_size: 8,
compaction_max_source_fragments: 32,
index_after_compaction: false,
compaction_max_bytes_per_file: 1024 * 1024 * 1024,
merge_wal_interval_secs: 0,
merge_wal_min_generations: 8,
Expand Down
Loading