From c03aef79421adebc3606a60a57622ee033e2b6b4 Mon Sep 17 00:00:00 2001 From: Beinan Date: Wed, 22 Jul 2026 17:40:57 +0000 Subject: [PATCH] perf(master): speed up rollout record browsing --- Cargo.lock | 1 + .../lance-context-core/src/rollout_store.rs | 252 +++++++++++++++++- crates/lance-context-master/Cargo.toml | 1 + crates/lance-context-master/src/routes.rs | 20 +- crates/lance-context-master/src/state.rs | 43 ++- 5 files changed, 301 insertions(+), 16 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 536b329..1acc3f9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5513,6 +5513,7 @@ dependencies = [ "lance-context-api", "lance-context-core", "lance-context-metrics", + "lru", "metrics", "reqwest 0.12.28", "serde", diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index 1bb9662..093f2b8 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -328,6 +328,16 @@ impl RolloutStore { Ok(()) } + /// Refresh this handle to the latest base-table manifest while retaining + /// its session and metadata caches. + /// + /// Long-lived read handles call this before a new request so compaction or + /// WAL merges committed by another process become visible without paying + /// the cost of reopening the dataset and rebuilding all session caches. + pub async fn refresh_latest(&mut self) -> LanceResult<()> { + self.dataset.checkout_latest().await + } + /// Append rollout rows through this instance's MemWAL shard; returns the /// current base dataset version. /// @@ -965,8 +975,16 @@ impl RolloutStore { /// Filter and page rollout rows in the LSM execution plan. /// /// Reads one row beyond the requested page to report `has_more`, avoiding - /// an unbounded full-table count on every UI request. Artifact bytes are - /// projected out, matching [`Self::list`]. + /// an unbounded full-table count on every UI request. Pagination is + /// deliberately late-materialized in two scans: + /// + /// 1. scan, sort, and deduplicate only `id` to select the page; + /// 2. fetch the complete non-blob columns for those page ids in one query. + /// + /// [`LsmScanner`] sorts every source by primary key before applying its + /// global limit. Keeping wide token/logprob/metadata columns out of that + /// full-source sort makes browsing large rollout tables substantially + /// cheaper while preserving the same LSM deduplication semantics. pub async fn list_filtered( &self, filters: &RolloutFilters, @@ -976,23 +994,49 @@ impl RolloutStore { let shard_snapshots = self.wal_shard_snapshots().await?; let filter = filters.expression(); - let columns = self.non_blob_columns(); - let refs: Vec<&str> = columns.iter().map(String::as_str).collect(); let mut page_scanner = self - .lsm_scanner_with_snapshots(shard_snapshots) - .project(&refs); + .lsm_scanner_with_snapshots(shard_snapshots.clone()) + .project(&["id"]); if let Some(filter) = &filter { page_scanner = page_scanner.filter(filter)?; } page_scanner = page_scanner.limit(limit.saturating_add(1), Some(offset)); let mut stream = page_scanner.try_into_stream().await?; - let mut records = Vec::new(); + let mut page_ids = Vec::new(); + while let Some(batch) = stream.try_next().await? { + let ids = column_as::(&batch, "id")?; + page_ids.extend((0..batch.num_rows()).map(|row| ids.value(row).to_string())); + } + let has_more = page_ids.len() > limit; + page_ids.truncate(limit); + if page_ids.is_empty() { + return Ok(RolloutPage { + records: Vec::new(), + has_more, + }); + } + + let columns = self.non_blob_columns(); + let refs: Vec<&str> = columns.iter().map(String::as_str).collect(); + let id_refs: Vec<&str> = page_ids.iter().map(String::as_str).collect(); + let id_filter = format!("id IN ({})", sql_quoted_list(&id_refs)); + let record_scanner = self + .lsm_scanner_with_snapshots(shard_snapshots) + .project(&refs) + .filter(&id_filter)?; + + let mut stream = record_scanner.try_into_stream().await?; + let mut records_by_id = HashMap::with_capacity(page_ids.len()); while let Some(batch) = stream.try_next().await? { - records.extend(batch_to_rollout_records(&batch)?); + for record in batch_to_rollout_records(&batch)? { + records_by_id.insert(record.id.clone(), record); + } } - let has_more = records.len() > limit; - records.truncate(limit); + let records = page_ids + .into_iter() + .filter_map(|id| records_by_id.remove(&id)) + .collect(); Ok(RolloutPage { records, has_more }) } @@ -1871,6 +1915,14 @@ fn optional_i8_list(array: Option<&ListArray>, row: usize) -> LanceResult