diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index b3f10f2..bc1b1b9 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -1233,6 +1233,11 @@ impl StorageBase { // master builds it before fan-out (#277), but a worker's own timer // or a manual merge must not depend on the master having been here. if self.dataset.count_fragments() > 0 && !self.has_key_btree_index().await? { + metrics::counter!("rollout_merge_index_built_on_demand_total").increment(1); + info!( + uri = %self.dataset.uri(), + "base table has no key BTree; building it before the merge" + ); self.create_key_btree_index().await?; } let key_index = merge_schema.index_of(&self.key_column)?; @@ -1747,7 +1752,7 @@ impl StorageBase { // the only cure is the merge that is already behind. // Fail the read fast so the caller retries after the // merge instead of taking the worker down with it. - metrics::counter!("rollout_reads_refused_pending_total").increment(1); + note_read_refused(&uri); return Err(LanceError::io(format!( "{PENDING_GENERATIONS_EXCEEDED}: shard {shard_id} has {pending} \ flushed generations pending merge (cap {max_at}); retry after \ @@ -1791,7 +1796,7 @@ impl StorageBase { .map(|snapshot| snapshot.flushed_generations.len()) .sum(); if max_at != 0 && total > max_at { - metrics::counter!("rollout_reads_refused_pending_total").increment(1); + note_read_refused(self.dataset.uri()); return Err(LanceError::io(format!( "{PENDING_GENERATIONS_EXCEEDED}: {total} flushed generations pending merge \ across {} shards (cap {max_at}); retry after the merge catches up", @@ -2198,6 +2203,21 @@ pub(crate) fn align_batch_to_schema( Ok(RecordBatch::try_new(target_schema, columns)?) } +/// A read refused because the store's WAL backlog is over the cap. Labeled by +/// store so the fleet can see *which* store is unreadable, not just that one +/// is: a single store answered 503 for five hours today before anyone +/// noticed. The label is the dataset directory name, bounded by the number +/// of stores a worker touches. +fn note_read_refused(uri: &str) { + let store = uri + .trim_end_matches('/') + .rsplit('/') + .next() + .unwrap_or("unknown") + .to_string(); + metrics::counter!("rollout_reads_refused_pending_total", "store" => store).increment(1); +} + /// Derive the MemWAL shard UUID a server instance writes to from its stable /// instance id. Deterministic (UUID v5), so the same instance id always maps to /// the same shard across restarts and reopens — the losing/gaining semantics of diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index 806d020..6336fb4 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -360,6 +360,7 @@ async fn scan_once_inner(state: &Arc) -> lance::Result { } } snapshot.sort_by(|a, b| a.name.cmp(&b.name)); + record_backlog_gauges(&snapshot); // Retire experiments with no writes for the configured window. Each is // merged, compacted and verified quiescent first; only then is its row @@ -521,6 +522,66 @@ async fn pending_or_timeout( } } +/// Fleet-wide view of the WAL backlog, refreshed every scan. The read cap +/// (`ROLLOUT_WAL_PENDING_MAX_GENERATIONS`, 4,096 by default) turns pending +/// generations into refused reads, so `stores_over_read_cap` is the number of +/// stores whose reads are failing right now; it is the alert that was +/// missing when a store sat at 17k for hours. +fn record_backlog_gauges(snapshot: &[StatRow]) { + let b = Backlog::summarize(snapshot); + metrics::gauge!("master_wal_pending_generations_total").set(b.total as f64); + metrics::gauge!("master_wal_pending_generations_max").set(b.max as f64); + metrics::gauge!("master_stores_pending_over_read_cap").set(b.over_read_cap as f64); + metrics::gauge!("master_stores_pending_over_1k").set(b.over_1k as f64); +} + +#[derive(Debug, Default, PartialEq, Eq)] +struct Backlog { + total: i64, + max: i64, + over_read_cap: i64, + over_1k: i64, +} + +impl Backlog { + const READ_CAP: i64 = lance_context_core::DEFAULT_PENDING_GENERATIONS_MAX as i64; + + fn summarize(snapshot: &[StatRow]) -> Self { + let mut b = Self::default(); + for row in snapshot { + let p = row.pending_wal_generations.max(0); + b.total += p; + b.max = b.max.max(p); + if p > Self::READ_CAP { + b.over_read_cap += 1; + } + if p > 1000 { + b.over_1k += 1; + } + } + b + } +} + +/// Pending generations that moved while the base-table version did not. +/// This is exactly the state that hid a 17k-generation backlog from the +/// merge sweep before #281: the version-unchanged shortcut reused the stale +/// count. Counted so it can be alerted on; the row itself is corrected. +fn note_pending_drift(name: &str, previous: i64, current: usize) { + let current = current as i64; + if current != previous { + metrics::counter!("master_stats_pending_recounted_total").increment(1); + if current > previous.max(0) + 256 { + tracing::warn!( + store = %name, + previous, + current, + "pending WAL generations grew while the base version stood still" + ); + } + } +} + /// Observe one generic store. Only the fields the WAL-merge sweep and the UI /// need: version (for the skip-if-unchanged shortcut), base row count, /// fragment count and pending MemWAL generations. Compaction counters stay at @@ -553,6 +614,7 @@ async fn observe_generic( if let Some(prev) = previous { if prev.version != StatRow::UNKNOWN_VERSION && prev.version == current_version { let pending = pending_or_timeout(name, store.pending_wal_generations()).await?; + note_pending_drift(name, prev.pending_wal_generations, pending); let mut row = prev.clone(); row.uri = uri.to_string(); row.pending_wal_generations = pending as i64; @@ -622,6 +684,7 @@ async fn observe_one( if let Some(prev) = previous { if prev.version != StatRow::UNKNOWN_VERSION && prev.version == current_version { let pending = pending_or_timeout(name, store.pending_wal_generations()).await?; + note_pending_drift(name, prev.pending_wal_generations, pending); let mut row = prev.clone(); row.uri = uri.to_string(); row.pending_wal_generations = pending as i64; @@ -1320,6 +1383,25 @@ mod retirement_tests { /// experiment is invisible to them forever. Dropping one with pending /// generations would strand them: never merged, read amplification never /// recovers, `_mem_wal/` never reclaimed, and no process left to notice. + #[test] + fn backlog_summary_counts_stores_over_the_read_cap() { + let mut a = row("a", "u", 0); + a.pending_wal_generations = 17_254; + let mut b = row("b", "u", 0); + b.pending_wal_generations = 1_200; + let mut c = row("c", "u", 0); + c.pending_wal_generations = 3; + assert_eq!( + Backlog::summarize(&[a, b, c]), + Backlog { + total: 18_457, + max: 17_254, + over_read_cap: 1, + over_1k: 2, + } + ); + } + #[tokio::test] async fn retirement_drains_the_wal_before_dropping_the_row() { let dir = TempDir::new().unwrap();