Skip to content
Merged
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
94 changes: 93 additions & 1 deletion crates/lance-context-master/src/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -507,6 +507,20 @@ struct ScanEntry {
kind: ScanKind,
}

/// Count pending generations under the observe timeout; shard manifests are
/// one small object read per shard, but object storage still has to answer.
async fn pending_or_timeout(
name: &str,
pending: impl std::future::Future<Output = lance::Result<usize>>,
) -> lance::Result<usize> {
match tokio::time::timeout(OBSERVE_TIMEOUT, pending).await {
Ok(result) => result,
Err(_) => Err(lance::Error::io(format!(
"pending-generation count timed out for store '{name}'"
))),
}
}

/// 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 All @@ -529,11 +543,19 @@ async fn observe_generic(
}
};

// WAL flushes never bump the base-table version, so the unchanged-
// version shortcut must still recount pending generations: a store that
// no worker or master has merged keeps its old version while its shards
// accumulate thousands of generations, and the merge sweep reads this
// row to decide whether to merge it at all. Only the row count (a
// base-table scan) is carried over.
let current_version = store.version() as i64;
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?;
let mut row = prev.clone();
row.uri = uri.to_string();
row.pending_wal_generations = pending as i64;
row.scanned_at = Utc::now().timestamp_millis();
return Ok((row, true));
}
Expand Down Expand Up @@ -593,12 +615,16 @@ async fn observe_one(
}
};

// Free: the manifest is already in hand from the open above.
// Free: the manifest is already in hand from the open above. Pending
// generations are not: see `observe_generic` for why they are recounted
// even when the base-table version is unchanged.
let current_version = store.version() as i64;
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?;
let mut row = prev.clone();
row.uri = uri.to_string();
row.pending_wal_generations = pending as i64;
row.scanned_at = Utc::now().timestamp_millis();
return Ok((row, true));
}
Expand Down Expand Up @@ -1053,6 +1079,72 @@ mod incremental_scan_tests {
.unwrap();
assert!(skipped);
assert_eq!(again.pending_wal_generations, 3);

// More flushes do not move the base version, but the merge sweep
// reads this count: the skipped row must still carry the new one.
{
let store = GenericStore::open(
&uri,
SchemaSpec::new(vec![(
ID_COLUMN.to_string(),
ColumnSpec::required(ColumnType::String { large: false }),
)]),
GenericStoreOptions {
seal_on_add: true,
..Default::default()
},
)
.await
.unwrap();
for i in 3..5 {
store
.add(&[serde_json::json!({"id": format!("r{i}")})
.as_object()
.unwrap()
.clone()])
.await
.unwrap();
}
}
let (recounted, skipped) =
observe_generic(&name, &uri, Some(&again), GenericStoreOptions::default())
.await
.unwrap();
assert!(skipped, "base version unchanged");
assert_eq!(recounted.version, again.version);
assert_eq!(recounted.pending_wal_generations, 5);
}

/// WAL flushes never bump the base-table version. A store that no one
/// has merged keeps its version forever while its shards fill up; the
/// unchanged-version shortcut must still recount pending generations or
/// the merge sweep never learns about it (17k generations sat behind a
/// stale "3" in production).
#[tokio::test]
async fn skipped_round_recounts_pending_generations() {
let dir = TempDir::new().unwrap();
let uri = dir.path().join("e.lance").to_string_lossy().to_string();
let store = RolloutStore::open_with_options(&uri, RolloutStoreOptions::default())
.await
.unwrap();
store.add(&[rec("a")]).await.unwrap();
store.flush().await.unwrap();

let (prev, _) = observe_one("e", &uri, None, RolloutStoreOptions::default())
.await
.unwrap();
assert_eq!(prev.pending_wal_generations, 1);

for id in ["b", "c", "d"] {
store.add(&[rec(id)]).await.unwrap();
store.flush().await.unwrap();
}
let (next, skipped) = observe_one("e", &uri, Some(&prev), RolloutStoreOptions::default())
.await
.unwrap();
assert!(skipped, "base version unchanged");
assert_eq!(next.version, prev.version);
assert_eq!(next.pending_wal_generations, 4);
}

/// A write moves the base version, so the next scan must observe fully and
Expand Down
Loading