diff --git a/crates/lance-context-master/src/config.rs b/crates/lance-context-master/src/config.rs index 753c684..bba4dcd 100644 --- a/crates/lance-context-master/src/config.rs +++ b/crates/lance-context-master/src/config.rs @@ -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, diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index d5a0a25..cf648ab 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -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, diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 1d5b197..0ae6cb0 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -143,7 +143,7 @@ async fn run_task(state: &Arc, 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, }; @@ -174,15 +174,39 @@ async fn run_task(state: &Arc, 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, target: &str) -> Result { - 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, task: &TaskRecord) -> Result { + 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 @@ -209,7 +233,10 @@ async fn index_id_inner(state: &Arc, name: &str) -> Result, name: &str) -> Result { +async fn compact_inner( + state: &Arc, + name: &str, +) -> Result { let uri = state.rollout_uri(name); let config = state.compaction_config(); let _permit = state @@ -227,10 +254,7 @@ async fn compact_inner(state: &Arc, name: &str) -> Result dispatcher builds the ZoneMap /// index -> task reaches Done with the expected detail summary. #[tokio::test] diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index dba767d..b3dadbf 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -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, diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index b522c90..d3e5b43 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -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,