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
24 changes: 22 additions & 2 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?;
Expand Down Expand Up @@ -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 \
Expand Down Expand Up @@ -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",
Expand Down Expand 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
Expand Down
82 changes: 82 additions & 0 deletions crates/lance-context-master/src/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -360,6 +360,7 @@ async fn scan_once_inner(state: &Arc<MasterState>) -> lance::Result<usize> {
}
}
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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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();
Expand Down
Loading