This is an automated email from the ASF dual-hosted git repository.
jerry-024 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new ab90d011 feat: add vector search timing and I/O diagnostics (#711)
ab90d011 is described below
commit ab90d01102a989298c7e03c5bae4feaecd9f5bb0
Author: jerry <[email protected]>
AuthorDate: Mon Aug 17 15:02:32 2026 +0800
feat: add vector search timing and I/O diagnostics (#711)
---
crates/paimon/src/io/cache/disk.rs | 7 +-
crates/paimon/src/table/vector_search_builder.rs | 392 +++++++++++++++++++++--
crates/paimon/src/vindex/mod.rs | 32 ++
crates/paimon/src/vindex/range_reader.rs | 178 +++++++++-
crates/paimon/src/vindex/reader.rs | 146 +++++++--
5 files changed, 697 insertions(+), 58 deletions(-)
diff --git a/crates/paimon/src/io/cache/disk.rs
b/crates/paimon/src/io/cache/disk.rs
index bd19e93a..771282ee 100644
--- a/crates/paimon/src/io/cache/disk.rs
+++ b/crates/paimon/src/io/cache/disk.rs
@@ -875,12 +875,9 @@ mod tests {
async fn test_disk_cache_restart_defers_crc_validation_until_first_hit() {
let directory = tempfile::tempdir().unwrap();
let key = BlockKey::new("s3://bucket/table/snapshot/snapshot-1", 4, 0);
- let cache = DiskCache::new(directory.path(), None).unwrap();
- cache.put_block(&key, Bytes::from_static(b"data")).await;
- drop(cache);
-
let block_path = directory.path().join(key.cache_relative_path());
- let mut encoded = std::fs::read(&block_path).unwrap();
+ std::fs::create_dir_all(block_path.parent().unwrap()).unwrap();
+ let mut encoded = encode_block(&key, &Bytes::from_static(b"data"));
let payload_offset = encoded.len() - CHECKSUM_LEN - 1;
encoded[payload_offset] ^= 0xff;
std::fs::write(&block_path, encoded).unwrap();
diff --git a/crates/paimon/src/table/vector_search_builder.rs
b/crates/paimon/src/table/vector_search_builder.rs
index f5352206..5919b4c7 100644
--- a/crates/paimon/src/table/vector_search_builder.rs
+++ b/crates/paimon/src/table/vector_search_builder.rs
@@ -58,9 +58,9 @@ use crate::vindex::pkvector::ann::{AnnSegmentSource,
PkVectorAnnSearcher, Vindex
use crate::vindex::pkvector::bucket::{BucketActiveFile, BucketAnnSegment,
ExactFileSearchFuture};
use crate::vindex::pkvector::exact::validate_query;
use crate::vindex::pkvector::metric::VectorSearchMetric;
-use crate::vindex::range_reader::VindexFileReader;
+use crate::vindex::range_reader::{RangeIoStats, VindexFileReader};
use crate::vindex::reader::VindexVectorGlobalIndexReader;
-use crate::vindex::{is_vindex_index_type, VindexVectorIndexOptions};
+use crate::vindex::{is_vindex_index_type, vector_search_timing_enabled,
VindexVectorIndexOptions};
use arrow_array::{Array, FixedSizeListArray, Float32Array, Int64Array,
ListArray, RecordBatch};
use arrow_select::interleave::interleave_record_batch;
use futures::{stream, TryStreamExt};
@@ -73,6 +73,7 @@ use std::cmp::Ordering;
use std::collections::{BinaryHeap, HashMap, HashSet};
use std::io::Cursor;
use std::sync::Arc;
+use std::time::{Duration, Instant};
const INDEX_DIR: &str = "index";
@@ -139,6 +140,23 @@ fn vindex_index_parallelism(entry_count: usize,
max_concurrency: usize) -> usize
entry_count.min(max_concurrency).max(1)
}
+fn log_vindex_range_io_stats(file: &str, query_count: usize, stats:
&RangeIoStats) {
+ let stats = stats.snapshot();
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_range_io file={} nq={} logical_ranges={}
requested_bytes={} file_read_calls={} returned_bytes={} read_ahead_hits={}
io_wait_sum_ms={:.3} range_permit_wait_sum_ms={:.3}",
+ file,
+ query_count,
+ stats.logical_ranges,
+ stats.requested_bytes,
+ stats.file_read_calls,
+ stats.returned_bytes,
+ stats.read_ahead_hits,
+ stats.io_wait_nanos as f64 / 1_000_000.0,
+ stats.range_permit_wait_nanos as f64 / 1_000_000.0,
+ );
+}
+
pub struct VectorSearchBuilder<'a> {
table: &'a Table,
vector_column: Option<String>,
@@ -920,9 +938,11 @@ async fn plan_and_search_pk_candidates_batch(
reader.visit_batch_vector_search(searches, |_|
Ok(Cursor::new(data)))
}
(VectorIndexBackend::Vindex, AnnSegmentSource::Vindex(source))
=> {
+ let range_io_stats = source.range_io_stats();
let mut reader =
VindexVectorGlobalIndexReader::new(io_meta, options.clone())
.with_batch_index_parallelism(batch_index_parallelism);
- reader.load_validated(
+ let results = reader.visit_batch_vector_search_validated(
+ searches,
|_| Ok(source),
|metadata| {
verify_segment_metric(
@@ -931,7 +951,10 @@ async fn plan_and_search_pk_candidates_batch(
)
},
)?;
- reader.search_batch(searches)
+ if let Some(stats) = range_io_stats {
+ log_vindex_range_io_stats(&segment.path,
searches.len(), &stats);
+ }
+ Ok(results)
}
(VectorIndexBackend::Lumina, AnnSegmentSource::Vindex(_))
| (VectorIndexBackend::Vindex, AnnSegmentSource::Buffered(_))
=> {
@@ -1167,6 +1190,8 @@ impl<'a> BatchVectorSearchBuilder<'a> {
}
pub async fn execute(&self) -> crate::Result<Vec<SearchResult>> {
+ let timing_enabled = vector_search_timing_enabled();
+ let total_start = timing_enabled.then(Instant::now);
// Fail closed: like `execute_read` and the single-query builder, this
// returns data-derived row ids/scores outside `TableScan`/`TableRead`,
// so it must refuse a `query-auth.enabled` table before any fast path
@@ -1242,12 +1267,33 @@ impl<'a> BatchVectorSearchBuilder<'a> {
.collect::<crate::Result<Vec<_>>>()?;
let snapshot_manager = self.table.snapshot_manager();
+ let setup = total_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let snapshot_start = timing_enabled.then(Instant::now);
let snapshot = match
crate::table::time_travel::resolve_snapshot(self.table).await? {
Some(s) => s,
- None => return Ok(vec![SearchResult::empty();
vector_searches.len()]),
+ None => {
+ let snapshot = snapshot_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let results = vec![SearchResult::empty();
vector_searches.len()];
+ if let Some(total_start) = total_start {
+ let total = total_start.elapsed();
+ let unattributed =
total.saturating_sub(setup.saturating_add(snapshot));
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_search_api nq={} index_entries=0
result_count=0 total_ms={:.3} setup_ms={:.3} snapshot_ms={:.3}
manifest_ms=0.000 evaluate_ms=0.000 unattributed_ms={:.3}",
+ vector_searches.len(),
+ total.as_secs_f64() * 1000.0,
+ setup.as_secs_f64() * 1000.0,
+ snapshot.as_secs_f64() * 1000.0,
+ unattributed.as_secs_f64() * 1000.0,
+ );
+ }
+ return Ok(results);
+ }
};
+ let snapshot_elapsed = snapshot_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let manifest_start = timing_enabled.then(Instant::now);
let index_entries = match snapshot.index_manifest() {
Some(index_manifest_name) => {
let manifest_path =
snapshot_manager.manifest_path(index_manifest_name);
@@ -1255,8 +1301,10 @@ impl<'a> BatchVectorSearchBuilder<'a> {
}
None => Vec::new(),
};
+ let manifest = manifest_start.map_or(Duration::ZERO, |start|
start.elapsed());
- evaluate_batch_vector_search(
+ let evaluate_start = timing_enabled.then(Instant::now);
+ let results = evaluate_batch_vector_search(
VectorSearchEvaluation {
table: Some(self.table),
file_io: self.table.file_io(),
@@ -1268,7 +1316,33 @@ impl<'a> BatchVectorSearchBuilder<'a> {
&index_entries,
&vector_searches,
)
- .await
+ .await?;
+ if let (Some(total_start), Some(evaluate_start)) = (total_start,
evaluate_start) {
+ let total = total_start.elapsed();
+ let evaluate = evaluate_start.elapsed();
+ let children = setup
+ .saturating_add(snapshot_elapsed)
+ .saturating_add(manifest)
+ .saturating_add(evaluate);
+ let result_count = results
+ .iter()
+ .map(|result| result.row_ids.len())
+ .sum::<usize>();
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_search_api nq={} index_entries={}
result_count={} total_ms={:.3} setup_ms={:.3} snapshot_ms={:.3}
manifest_ms={:.3} evaluate_ms={:.3} unattributed_ms={:.3}",
+ vector_searches.len(),
+ index_entries.len(),
+ result_count,
+ total.as_secs_f64() * 1000.0,
+ setup.as_secs_f64() * 1000.0,
+ snapshot_elapsed.as_secs_f64() * 1000.0,
+ manifest.as_secs_f64() * 1000.0,
+ evaluate.as_secs_f64() * 1000.0,
+ total.saturating_sub(children).as_secs_f64() * 1000.0,
+ );
+ }
+ Ok(results)
}
/// Run a batch of vector searches and materialize each query's matching
rows as
@@ -1423,6 +1497,12 @@ struct VectorSearchEvaluation<'a> {
next_row_id: Option<i64>,
}
+#[derive(Default)]
+struct IndexSearchTiming {
+ permit_wait: Duration,
+ file_reader_open: Duration,
+}
+
#[cfg(test)]
async fn evaluate_vector_search(
evaluation: VectorSearchEvaluation<'_>,
@@ -1444,6 +1524,8 @@ async fn evaluate_batch_vector_search(
index_entries: &[IndexManifestEntry],
vector_searches: &[VectorSearch],
) -> crate::Result<Vec<SearchResult>> {
+ let timing_enabled = vector_search_timing_enabled();
+ let total_start = timing_enabled.then(Instant::now);
if vector_searches.is_empty() {
return Ok(Vec::new());
}
@@ -1495,6 +1577,7 @@ async fn evaluate_batch_vector_search(
return Ok(vec![SearchResult::empty(); vector_searches.len()]);
}
+ let deletion_vector_start = timing_enabled.then(Instant::now);
let deleted_row_index = if core_options.data_evolution_enabled() {
match evaluation.table {
Some(table) => {
@@ -1507,6 +1590,7 @@ async fn evaluate_batch_vector_search(
} else {
None
};
+ let deletion_vector = deletion_vector_start.map_or(Duration::ZERO, |start|
start.elapsed());
let max_limit = vector_searches
.iter()
@@ -1524,8 +1608,16 @@ async fn evaluate_batch_vector_search(
};
let index_search_limit = indexed_search_limit(max_limit, refine_factor)?;
+ let vector_entry_count = vector_entries.len();
+ let mut permit_wait = Duration::ZERO;
+ let mut file_reader_open = Duration::ZERO;
+ let mut index_search = Duration::ZERO;
+ let mut merge = Duration::ZERO;
+ let mut refine = Duration::ZERO;
+ let mut raw_fallback = Duration::ZERO;
let mut merged = vec![SearchResult::empty(); vector_searches.len()];
if !vector_entries.is_empty() {
+ let index_search_start = timing_enabled.then(Instant::now);
let concurrency = core_options.global_index_thread_num()?;
if concurrency > tokio::sync::Semaphore::MAX_PERMITS {
return Err(crate::Error::DataInvalid {
@@ -1573,13 +1665,19 @@ async fn evaluate_batch_vector_search(
options.extend(search_options.clone());
let input = evaluation.file_io.new_input(&path);
async move {
+ let permit_start = timing_enabled.then(Instant::now);
let permit =
acquire_process_global_search_permit(concurrency).await?;
+ let permit_wait =
+ permit_start.map_or(Duration::ZERO, |start|
start.elapsed());
let input = input?;
let query_count = vector_searches.len();
+ let mut file_reader_open = Duration::ZERO;
+ let mut full_file_read = None;
let io_meta =
GlobalIndexIOMeta::new(file_name.clone(), file_size,
index_meta_bytes);
let results = match backend {
VectorIndexBackend::Lumina => {
+ let read_start = timing_enabled.then(Instant::now);
let data = input.read().await.map_err(|e| {
crate::Error::DataInvalid {
message: format!(
@@ -1591,6 +1689,9 @@ async fn evaluate_batch_vector_search(
source: None,
}
})?;
+ if let Some(start) = read_start {
+ full_file_read = Some((start.elapsed(),
data.len()));
+ }
execute_global_index_with_guard(
"Lumina global-index batch search task failed",
permit,
@@ -1607,6 +1708,8 @@ async fn evaluate_batch_vector_search(
VectorIndexBackend::Vindex => {
match tokio::runtime::Handle::try_current() {
Ok(runtime) => {
+ let file_reader_open_start =
+ timing_enabled.then(Instant::now);
let file_reader =
input.reader().await.map_err(|e| {
crate::Error::DataInvalid {
message: format!(
@@ -1616,6 +1719,8 @@ async fn evaluate_batch_vector_search(
source: None,
}
})?;
+ file_reader_open = file_reader_open_start
+ .map_or(Duration::ZERO, |start|
start.elapsed());
let source =
VindexFileReader::new_with_permits(
Arc::new(file_reader),
runtime,
@@ -1623,18 +1728,28 @@ async fn evaluate_batch_vector_search(
file_size,
file_name.clone(),
);
- execute_vindex_searches(
+ let range_io_stats =
source.range_io_stats();
+ let results = execute_vindex_searches(
io_meta,
options,
vector_searches,
source,
- file_name,
+ file_name.clone(),
batch_index_parallelism,
permit,
)
- .await?
+ .await?;
+ if let Some(stats) = range_io_stats {
+ log_vindex_range_io_stats(
+ &file_name,
+ query_count,
+ &stats,
+ );
+ }
+ results
}
Err(_) if query_count > 1 => {
+ let read_start =
timing_enabled.then(Instant::now);
let data = input.read().await.map_err(|e| {
crate::Error::DataInvalid {
message: format!(
@@ -1644,12 +1759,15 @@ async fn evaluate_batch_vector_search(
source: None,
}
})?;
+ if let Some(start) = read_start {
+ full_file_read =
Some((start.elapsed(), data.len()));
+ }
execute_vindex_searches(
io_meta,
options,
vector_searches,
Cursor::new(data),
- file_name,
+ file_name.clone(),
batch_index_parallelism,
permit,
)
@@ -1666,6 +1784,18 @@ async fn evaluate_batch_vector_search(
}
}
};
+ if let Some((read, returned_bytes)) = full_file_read {
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_full_file_io backend={}
file={} nq={} requested_bytes={} returned_bytes={} read_ms={:.3}",
+ backend.error_name(),
+ file_name,
+ query_count,
+ file_size,
+ returned_bytes,
+ read.as_secs_f64() * 1000.0,
+ );
+ }
if results.len() != query_count {
return Err(crate::Error::DataInvalid {
message: format!(
@@ -1677,7 +1807,7 @@ async fn evaluate_batch_vector_search(
});
}
- Ok::<_, crate::Error>(
+ Ok::<_, crate::Error>((
results
.into_iter()
.map(|result| match result {
@@ -1686,20 +1816,30 @@ async fn evaluate_batch_vector_search(
None => SearchResult::empty(),
})
.collect::<Vec<_>>(),
- )
+ IndexSearchTiming {
+ permit_wait,
+ file_reader_open,
+ },
+ ))
}
})
.collect();
let results = drain_indexed_jobs(futures.into_iter(),
concurrency).await?;
- for per_entry in &results {
+ index_search = index_search_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let merge_start = timing_enabled.then(Instant::now);
+ for (per_entry, entry_timing) in &results {
+ permit_wait = permit_wait.saturating_add(entry_timing.permit_wait);
+ file_reader_open =
file_reader_open.saturating_add(entry_timing.file_reader_open);
for (query_index, result) in per_entry.iter().enumerate() {
merged[query_index] = merged[query_index].or(result);
}
}
+ merge = merge_start.map_or(Duration::ZERO, |start| start.elapsed());
}
if refine_factor != 0 {
+ let refine_start = timing_enabled.then(Instant::now);
merged = maybe_rerank_indexed_batch_results(
evaluation,
index_entries,
@@ -1710,9 +1850,11 @@ async fn evaluate_batch_vector_search(
index_search_limit,
)
.await?;
+ refine = refine_start.map_or(Duration::ZERO, |start| start.elapsed());
}
if search_mode != GlobalIndexSearchMode::Fast {
+ let raw_fallback_start = timing_enabled.then(Instant::now);
let detail_ranges = if search_mode == GlobalIndexSearchMode::Detail {
let table = evaluation.table.ok_or_else(||
crate::Error::DataInvalid {
message: "Vector raw search in detail mode requires table
context".to_string(),
@@ -1736,6 +1878,7 @@ async fn evaluate_batch_vector_search(
message: "Vector raw search requires table
context".to_string(),
source: None,
})?;
+ let metric_start = timing_enabled.then(Instant::now);
let metric = resolve_raw_vector_metric(
evaluation.file_io,
table_path,
@@ -1745,15 +1888,35 @@ async fn evaluate_batch_vector_search(
field_name,
)
.await?;
- let raw_results =
+ let metric_resolve = metric_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let (raw_results, raw_timing) =
read_raw_batch_vector_search(table, vector_searches,
&raw_ranges, metric).await?;
+ if let Some(raw_timing) = raw_timing {
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_raw_fallback nq={} row_ranges={}
metric_resolve_ms={:.3} raw_plan_ms={:.3} split_count={} file_count={}
raw_stream_wait_ms={:.3} raw_score_cpu_ms={:.3} arrow_batches={} arrow_rows={}
total_raw_read_ms={:.3}",
+ vector_searches.len(),
+ raw_ranges.len(),
+ metric_resolve.as_secs_f64() * 1000.0,
+ raw_timing.plan.as_secs_f64() * 1000.0,
+ raw_timing.split_count,
+ raw_timing.file_count,
+ raw_timing.stream_wait.as_secs_f64() * 1000.0,
+ raw_timing.score_cpu.as_secs_f64() * 1000.0,
+ raw_timing.batch_count,
+ raw_timing.row_count,
+ raw_timing.total.as_secs_f64() * 1000.0,
+ );
+ }
for (query_index, result) in raw_results.iter().enumerate() {
merged[query_index] = merged[query_index].or(result);
}
}
+ raw_fallback = raw_fallback_start.map_or(Duration::ZERO, |start|
start.elapsed());
}
- merged
+ let finalize_start = timing_enabled.then(Instant::now);
+ let results = merged
.into_iter()
.zip(vector_searches)
.map(|(result, vector_search)| {
@@ -1761,7 +1924,41 @@ async fn evaluate_batch_vector_search(
.without_deleted_row_ranges(deleted_row_index.as_ref())?
.top_k(vector_search.limit))
})
- .collect()
+ .collect::<crate::Result<Vec<_>>>()?;
+ let finalize = finalize_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ if let Some(total_start) = total_start {
+ let total = total_start.elapsed();
+ let children = deletion_vector
+ .saturating_add(index_search)
+ .saturating_add(merge)
+ .saturating_add(refine)
+ .saturating_add(raw_fallback)
+ .saturating_add(finalize);
+ let result_count = results
+ .iter()
+ .map(|result| result.row_ids.len())
+ .sum::<usize>();
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_search_evaluate nq={} index_entries={}
index_files={} result_count={} refine_factor={} total_ms={:.3}
deletion_vector_ms={:.3} index_search_ms={:.3} global_permit_wait_sum_ms={:.3}
file_reader_open_sum_ms={:.3} merge_ms={:.3} refine_ms={:.3}
raw_fallback_ms={:.3} finalize_ms={:.3} unattributed_ms={:.3}",
+ vector_searches.len(),
+ index_entries.len(),
+ vector_entry_count,
+ result_count,
+ refine_factor,
+ total.as_secs_f64() * 1000.0,
+ deletion_vector.as_secs_f64() * 1000.0,
+ index_search.as_secs_f64() * 1000.0,
+ permit_wait.as_secs_f64() * 1000.0,
+ file_reader_open.as_secs_f64() * 1000.0,
+ merge.as_secs_f64() * 1000.0,
+ refine.as_secs_f64() * 1000.0,
+ raw_fallback.as_secs_f64() * 1000.0,
+ finalize.as_secs_f64() * 1000.0,
+ total.saturating_sub(children).as_secs_f64() * 1000.0,
+ );
+ }
+ Ok(results)
}
fn is_vector_global_index_file(index_file: &IndexFileMeta) -> bool {
@@ -2314,12 +2511,16 @@ async fn maybe_rerank_indexed_batch_results(
results: Vec<SearchResult>,
index_search_limit: usize,
) -> crate::Result<Vec<SearchResult>> {
+ let timing_enabled = vector_search_timing_enabled();
+ let total_start = timing_enabled.then(Instant::now);
let mut candidate_searches = Vec::with_capacity(vector_searches.len());
let mut candidate_results = Vec::with_capacity(vector_searches.len());
let mut union_candidates = RoaringTreemap::new();
+ let mut candidate_references = 0usize;
for (result, vector_search) in results.into_iter().zip(vector_searches) {
let candidates = result.top_k(index_search_limit);
+ candidate_references =
candidate_references.saturating_add(candidates.row_ids.len());
let mut include_row_ids = RoaringTreemap::new();
for &row_id in &candidates.row_ids {
include_row_ids.insert(row_id);
@@ -2340,7 +2541,9 @@ async fn maybe_rerank_indexed_batch_results(
message: "Vector index rerank requires table context".to_string(),
source: None,
})?;
+ let unique_candidates = union_candidates.len();
let raw_ranges = sorted_row_ids_to_row_ranges(union_candidates.iter())?;
+ let metric_start = timing_enabled.then(Instant::now);
let metric = resolve_raw_vector_metric(
evaluation.file_io,
evaluation.table_path.trim_end_matches('/'),
@@ -2350,8 +2553,30 @@ async fn maybe_rerank_indexed_batch_results(
field_name,
)
.await?;
-
- read_raw_batch_vector_search(table, &candidate_searches, &raw_ranges,
metric).await
+ let metric_resolve = metric_start.map_or(Duration::ZERO, |start|
start.elapsed());
+
+ let (results, raw_timing) =
+ read_raw_batch_vector_search(table, &candidate_searches, &raw_ranges,
metric).await?;
+ if let (Some(total_start), Some(raw_timing)) = (total_start, raw_timing) {
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vector_refine nq={} candidate_references={}
unique_candidates={} row_ranges={} metric_resolve_ms={:.3} raw_plan_ms={:.3}
split_count={} file_count={} raw_stream_wait_ms={:.3} raw_score_cpu_ms={:.3}
arrow_batches={} arrow_rows={} total_refine_ms={:.3}",
+ vector_searches.len(),
+ candidate_references,
+ unique_candidates,
+ raw_ranges.len(),
+ metric_resolve.as_secs_f64() * 1000.0,
+ raw_timing.plan.as_secs_f64() * 1000.0,
+ raw_timing.split_count,
+ raw_timing.file_count,
+ raw_timing.stream_wait.as_secs_f64() * 1000.0,
+ raw_timing.score_cpu.as_secs_f64() * 1000.0,
+ raw_timing.batch_count,
+ raw_timing.row_count,
+ total_start.elapsed().as_secs_f64() * 1000.0,
+ );
+ }
+ Ok(results)
}
fn sorted_row_ids_to_row_ranges(
@@ -2645,17 +2870,31 @@ fn configured_raw_vector_metric(
Ok(inferred.unwrap_or(RawVectorMetric::L2))
}
+#[derive(Default)]
+struct RawVectorReadTiming {
+ plan: Duration,
+ stream_wait: Duration,
+ score_cpu: Duration,
+ total: Duration,
+ split_count: usize,
+ file_count: usize,
+ batch_count: usize,
+ row_count: usize,
+}
+
async fn read_raw_batch_vector_search(
table: &Table,
vector_searches: &[VectorSearch],
raw_ranges: &[RowRange],
metric: RawVectorMetric,
-) -> crate::Result<Vec<SearchResult>> {
+) -> crate::Result<(Vec<SearchResult>, Option<RawVectorReadTiming>)> {
+ let timing_enabled = vector_search_timing_enabled();
+ let total_start = timing_enabled.then(Instant::now);
if vector_searches.is_empty() {
- return Ok(Vec::new());
+ return Ok((Vec::new(), None));
}
if raw_ranges.is_empty() {
- return Ok(vec![SearchResult::empty(); vector_searches.len()]);
+ return Ok((vec![SearchResult::empty(); vector_searches.len()], None));
}
let field_name = &vector_searches[0].field_name;
@@ -2670,13 +2909,28 @@ async fn read_raw_batch_vector_search(
});
}
+ let plan_start = timing_enabled.then(Instant::now);
let mut read_builder = table.new_read_builder();
read_builder
.with_projection(&[field_name.as_str(), ROW_ID_FIELD_NAME])?
.with_row_ranges(raw_ranges.to_vec());
let plan = read_builder.new_scan().plan().await?;
+ let plan_elapsed = plan_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let split_count = plan.splits().len();
+ let file_count = plan
+ .splits()
+ .iter()
+ .map(|split| split.data_files().len())
+ .sum();
if plan.splits().is_empty() {
- return Ok(vec![SearchResult::empty(); vector_searches.len()]);
+ return Ok((
+ vec![SearchResult::empty(); vector_searches.len()],
+ total_start.map(|start| RawVectorReadTiming {
+ plan: plan_elapsed,
+ total: start.elapsed(),
+ ..RawVectorReadTiming::default()
+ }),
+ ));
}
let read = read_builder.new_read()?;
let mut stream = read.to_arrow(plan.splits())?;
@@ -2686,14 +2940,44 @@ async fn read_raw_batch_vector_search(
.iter()
.map(|vector_search| RawScoreTopK::new(vector_search.limit))
.collect::<Vec<_>>();
- while let Some(batch) = stream.try_next().await? {
+ let mut timing = timing_enabled.then(|| RawVectorReadTiming {
+ plan: plan_elapsed,
+ split_count,
+ file_count,
+ ..RawVectorReadTiming::default()
+ });
+ loop {
+ let stream_wait_start = timing_enabled.then(Instant::now);
+ let batch = stream.try_next().await?;
+ if let (Some(timing), Some(stream_wait_start)) = (&mut timing,
stream_wait_start) {
+ timing.stream_wait = timing
+ .stream_wait
+ .saturating_add(stream_wait_start.elapsed());
+ }
+ let Some(batch) = batch else {
+ break;
+ };
+ if let Some(timing) = &mut timing {
+ timing.batch_count += 1;
+ timing.row_count =
timing.row_count.saturating_add(batch.num_rows());
+ }
+ let score_start = timing_enabled.then(Instant::now);
collect_raw_batch_vector_batch(&batch, vector_searches, metric,
&scoring_plan, &mut top_k)?;
+ if let (Some(timing), Some(score_start)) = (&mut timing, score_start) {
+ timing.score_cpu =
timing.score_cpu.saturating_add(score_start.elapsed());
+ }
}
- Ok(top_k
- .into_iter()
- .map(RawScoreTopK::into_search_result)
- .collect())
+ if let (Some(timing), Some(total_start)) = (&mut timing, total_start) {
+ timing.total = total_start.elapsed();
+ }
+ Ok((
+ top_k
+ .into_iter()
+ .map(RawScoreTopK::into_search_result)
+ .collect(),
+ timing,
+ ))
}
struct RawScoringPlan {
@@ -3125,7 +3409,37 @@ mod tests {
use arrow_array::ArrayRef;
use arrow_array::Int32Array;
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema
as ArrowSchema};
- use std::sync::Arc;
+ use std::sync::{Arc, Mutex, Once};
+
+ const VECTOR_SEARCH_LOG_TARGET: &str = "paimon::vector_search";
+ static VECTOR_SEARCH_TEST_LOGGER: VectorSearchTestLogger =
VectorSearchTestLogger;
+ static VECTOR_SEARCH_TEST_LOGS: Mutex<Vec<String>> =
Mutex::new(Vec::new());
+
+ struct VectorSearchTestLogger;
+
+ impl log::Log for VectorSearchTestLogger {
+ fn enabled(&self, metadata: &log::Metadata<'_>) -> bool {
+ metadata.target() == VECTOR_SEARCH_LOG_TARGET && metadata.level()
<= log::Level::Debug
+ }
+
+ fn log(&self, record: &log::Record<'_>) {
+ if self.enabled(record.metadata()) {
+ VECTOR_SEARCH_TEST_LOGS
+ .lock()
+ .unwrap()
+ .push(record.args().to_string());
+ }
+ }
+
+ fn flush(&self) {}
+ }
+
+ fn reset_vector_search_test_logs() {
+ static INIT: Once = Once::new();
+ INIT.call_once(||
log::set_logger(&VECTOR_SEARCH_TEST_LOGGER).unwrap());
+ log::set_max_level(log::LevelFilter::Debug);
+ VECTOR_SEARCH_TEST_LOGS.lock().unwrap().clear();
+ }
fn l2_score(distance: f32) -> f32 {
VectorSearchMetric::L2.distance_to_score(distance)
@@ -4958,7 +5272,9 @@ mod tests {
// ---- search_pk_route: candidate-only producer returns candidates +
context ----
#[tokio::test]
- async fn search_pk_route_returns_candidates_and_source_context() {
+ async fn search_pk_route_returns_candidates_and_publishes_diagnostics() {
+ reset_vector_search_test_logs();
+ let _timing = crate::vindex::enable_vector_search_timing_for_test();
// query [0,1]: squared-L2 distances pos1=0 < pos2=1 < pos0=2, so the
// strict-gap top-2 is [pos1, pos2] (best-first, not physical order).
let table = build_committed_pk_vector_table(&[[1.0, 0.0], [0.0, 1.0],
[1.0, 1.0]]).await;
@@ -4994,6 +5310,22 @@ mod tests {
.all(|c| c.split_index < route.splits.len()),
"candidate split_index must refer into the returned splits"
);
+
+ let logs = VECTOR_SEARCH_TEST_LOGS.lock().unwrap();
+ assert!(
+ logs.iter().any(|entry| {
+ entry.contains("event=paimon_vindex_reader")
+ && entry.contains("vector-ivf-flat-route.index")
+ }),
+ "PK vector search must publish vindex reader timing"
+ );
+ assert!(
+ logs.iter().any(|entry| {
+ entry.contains("event=paimon_vector_range_io")
+ && entry.contains("vector-ivf-flat-route.index")
+ }),
+ "PK vector search must publish range-I/O timing"
+ );
}
/// A table with no snapshot at all (never written) yields empty candidates
diff --git a/crates/paimon/src/vindex/mod.rs b/crates/paimon/src/vindex/mod.rs
index 5b51ea21..da89d257 100644
--- a/crates/paimon/src/vindex/mod.rs
+++ b/crates/paimon/src/vindex/mod.rs
@@ -24,6 +24,9 @@ pub mod pkvector;
use crate::spec::{DataField, DataType};
use paimon_vindex_core::index::VectorIndexConfig;
use std::collections::HashMap;
+#[cfg(test)]
+use std::sync::atomic::{AtomicUsize, Ordering};
+use std::sync::OnceLock;
pub const IVF_FLAT_IDENTIFIER: &str = "ivf-flat";
pub const IVF_PQ_IDENTIFIER: &str = "ivf-pq";
@@ -34,6 +37,35 @@ const DEFAULT_NLIST: &str = "256";
const DEFAULT_PQ_M: &str = "16";
const DEFAULT_PQ_USE_OPQ: &str = "false";
const DEFAULT_TRAIN_SAMPLE_RATIO: f64 = 1.0;
+const VECTOR_SEARCH_TIMING_ENV: &str = "PAIMON_LOG_VECTOR_SEARCH_TIMING";
+
+#[cfg(test)]
+static VECTOR_SEARCH_TIMING_TEST_GUARDS: AtomicUsize = AtomicUsize::new(0);
+
+#[cfg(test)]
+pub(crate) struct VectorSearchTimingTestGuard;
+
+#[cfg(test)]
+impl Drop for VectorSearchTimingTestGuard {
+ fn drop(&mut self) {
+ VECTOR_SEARCH_TIMING_TEST_GUARDS.fetch_sub(1, Ordering::Relaxed);
+ }
+}
+
+#[cfg(test)]
+pub(crate) fn enable_vector_search_timing_for_test() ->
VectorSearchTimingTestGuard {
+ VECTOR_SEARCH_TIMING_TEST_GUARDS.fetch_add(1, Ordering::Relaxed);
+ VectorSearchTimingTestGuard
+}
+
+pub(crate) fn vector_search_timing_enabled() -> bool {
+ #[cfg(test)]
+ if VECTOR_SEARCH_TIMING_TEST_GUARDS.load(Ordering::Relaxed) > 0 {
+ return true;
+ }
+ static ENABLED: OnceLock<bool> = OnceLock::new();
+ *ENABLED.get_or_init(||
std::env::var_os(VECTOR_SEARCH_TIMING_ENV).is_some_and(|v| v == "1"))
+}
pub fn is_vindex_index_type(index_type: &str) -> bool {
matches!(index_type, IVF_FLAT_IDENTIFIER | IVF_PQ_IDENTIFIER)
diff --git a/crates/paimon/src/vindex/range_reader.rs
b/crates/paimon/src/vindex/range_reader.rs
index d9c5e158..0f5d965c 100644
--- a/crates/paimon/src/vindex/range_reader.rs
+++ b/crates/paimon/src/vindex/range_reader.rs
@@ -16,12 +16,15 @@
// under the License.
use crate::io::FileRead;
+use crate::vindex::vector_search_timing_enabled;
use bytes::Bytes;
use futures::future::try_join_all;
use paimon_vindex_core::io::{ReadRequest, SeekRead, SeekReadCapabilities};
use std::io;
use std::ops::Range;
+use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{mpsc, Arc};
+use std::time::Instant;
const SCALAR_READ_MAX: usize = 64;
const SCALAR_READ_AHEAD: u64 = 64 * 1024;
@@ -54,6 +57,42 @@ struct MergedRange {
requested_bytes: u64,
}
+#[derive(Debug, Default)]
+pub(crate) struct RangeIoStats {
+ logical_ranges: AtomicU64,
+ requested_bytes: AtomicU64,
+ file_read_calls: AtomicU64,
+ returned_bytes: AtomicU64,
+ read_ahead_hits: AtomicU64,
+ io_wait_nanos: AtomicU64,
+ range_permit_wait_nanos: AtomicU64,
+}
+
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+pub(crate) struct RangeIoStatsSnapshot {
+ pub(crate) logical_ranges: u64,
+ pub(crate) requested_bytes: u64,
+ pub(crate) file_read_calls: u64,
+ pub(crate) returned_bytes: u64,
+ pub(crate) read_ahead_hits: u64,
+ pub(crate) io_wait_nanos: u64,
+ pub(crate) range_permit_wait_nanos: u64,
+}
+
+impl RangeIoStats {
+ pub(crate) fn snapshot(&self) -> RangeIoStatsSnapshot {
+ RangeIoStatsSnapshot {
+ logical_ranges: self.logical_ranges.load(Ordering::Relaxed),
+ requested_bytes: self.requested_bytes.load(Ordering::Relaxed),
+ file_read_calls: self.file_read_calls.load(Ordering::Relaxed),
+ returned_bytes: self.returned_bytes.load(Ordering::Relaxed),
+ read_ahead_hits: self.read_ahead_hits.load(Ordering::Relaxed),
+ io_wait_nanos: self.io_wait_nanos.load(Ordering::Relaxed),
+ range_permit_wait_nanos:
self.range_permit_wait_nanos.load(Ordering::Relaxed),
+ }
+ }
+}
+
/// Bridges vindex-core's synchronous positional reads to Paimon's asynchronous
/// range reader. This type is consumed from a blocking search task; it
captures
/// the surrounding Tokio runtime so remote storage reads still run
asynchronously.
@@ -64,6 +103,7 @@ pub(crate) struct VindexFileReader {
file_size: u64,
path: String,
scalar_cache: Option<CachedRange>,
+ stats: Option<Arc<RangeIoStats>>,
}
impl VindexFileReader {
@@ -97,9 +137,14 @@ impl VindexFileReader {
file_size,
path,
scalar_cache: None,
+ stats: vector_search_timing_enabled().then(||
Arc::new(RangeIoStats::default())),
}
}
+ pub(crate) fn range_io_stats(&self) -> Option<Arc<RangeIoStats>> {
+ self.stats.clone()
+ }
+
fn validate_range(&self, pos: u64, len: usize) -> io::Result<Range<u64>> {
let end = pos.checked_add(len as u64).ok_or_else(|| {
io::Error::new(
@@ -127,6 +172,9 @@ impl VindexFileReader {
if buf.len() <= SCALAR_READ_MAX {
if let Some(cache) = &self.scalar_cache {
if cache.contains(&range) {
+ if let Some(stats) = &self.stats {
+ stats.read_ahead_hits.fetch_add(1, Ordering::Relaxed);
+ }
let start = (range.start - cache.start) as usize;
buf.copy_from_slice(&cache.data[start..start + buf.len()]);
return Ok(());
@@ -163,17 +211,30 @@ impl VindexFileReader {
let permits = Arc::clone(&self.permits);
let path = self.path.clone();
let requested = ranges.to_vec();
+ let stats = self.stats.clone();
let (sender, receiver) = mpsc::sync_channel(1);
+ let wait_start = self.stats.as_ref().map(|_| Instant::now());
self.runtime.spawn(async move {
let fetched = try_join_all(requested.iter().cloned().map(|range| {
let reader = Arc::clone(&reader);
let permits = Arc::clone(&permits);
let path = path.clone();
+ let stats = stats.clone();
async move {
- let _permit = permits.acquire_owned().await.map_err(|_| {
+ let permit_wait_start = stats.as_ref().map(|_|
Instant::now());
+ let permit = permits.acquire_owned().await;
+ if let (Some(stats), Some(start)) = (&stats,
permit_wait_start) {
+ stats
+ .range_permit_wait_nanos
+ .fetch_add(start.elapsed().as_nanos() as u64,
Ordering::Relaxed);
+ }
+ let _permit = permit.map_err(|_| {
io::Error::other("vindex range read concurrency
limiter closed")
})?;
let expected = (range.end - range.start) as usize;
+ if let Some(stats) = &stats {
+ stats.file_read_calls.fetch_add(1, Ordering::Relaxed);
+ }
let data =
reader.read(range.clone()).await.map_err(|error| {
io::Error::other(format!(
"failed to read vindex file '{path}' range {}..{}:
{error}",
@@ -191,13 +252,24 @@ impl VindexFileReader {
),
));
}
+ if let Some(stats) = &stats {
+ stats
+ .returned_bytes
+ .fetch_add(data.len() as u64, Ordering::Relaxed);
+ }
Ok(data)
}
}))
.await;
let _ = sender.send(fetched);
});
- receiver.recv().map_err(|_| {
+ let result = receiver.recv();
+ if let (Some(stats), Some(start)) = (&self.stats, wait_start) {
+ stats
+ .io_wait_nanos
+ .fetch_add(start.elapsed().as_nanos() as u64,
Ordering::Relaxed);
+ }
+ result.map_err(|_| {
io::Error::other(format!(
"vindex range read task for '{}' was cancelled",
self.path
@@ -266,6 +338,20 @@ impl VindexFileReader {
impl SeekRead for VindexFileReader {
fn pread(&mut self, requests: &mut [ReadRequest<'_>]) -> io::Result<()> {
+ if let Some(stats) = &self.stats {
+ let (logical_ranges, requested_bytes) = requests
+ .iter()
+ .filter(|request| !request.buf.is_empty())
+ .fold((0u64, 0u64), |(ranges, bytes), request| {
+ (ranges + 1, bytes.saturating_add(request.buf.len() as
u64))
+ });
+ stats
+ .logical_ranges
+ .fetch_add(logical_ranges, Ordering::Relaxed);
+ stats
+ .requested_bytes
+ .fetch_add(requested_bytes, Ordering::Relaxed);
+ }
let non_empty = requests
.iter()
.filter(|request| !request.buf.is_empty())
@@ -289,6 +375,7 @@ impl SeekRead for VindexFileReader {
file_size: self.file_size,
path: self.path.clone(),
scalar_cache: None,
+ stats: self.stats.clone(),
}))
}
@@ -592,6 +679,83 @@ mod tests {
assert_eq!(cloned.read_capabilities(),
SeekReadCapabilities::default());
}
+ #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+ async fn range_io_stats_are_shared_across_clones() {
+ let data = Bytes::from(vec![7u8; 1024]);
+ let source: Arc<dyn FileRead> = TrackingRead::new(data.clone());
+ let mut reader = VindexFileReader::new(
+ source,
+ tokio::runtime::Handle::current(),
+ data.len() as u64,
+ "index".to_string(),
+ );
+ let stats = Arc::new(RangeIoStats::default());
+ reader.stats = Some(Arc::clone(&stats));
+ let mut cloned = reader.try_clone_reader().unwrap().unwrap();
+
+ tokio::task::spawn_blocking(move || {
+ let mut first = [0u8; 128];
+ reader
+ .pread(&mut [ReadRequest::new(0, &mut first)])
+ .unwrap();
+ let mut second = [0u8; 128];
+ cloned
+ .pread(&mut [ReadRequest::new(128, &mut second)])
+ .unwrap();
+ })
+ .await
+ .unwrap();
+
+ let stats = stats.snapshot();
+ assert_eq!(
+ (
+ stats.logical_ranges,
+ stats.requested_bytes,
+ stats.file_read_calls,
+ stats.returned_bytes,
+ stats.read_ahead_hits,
+ ),
+ (2, 256, 2, 256, 0)
+ );
+ assert!(stats.io_wait_nanos > 0);
+ }
+
+ #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+ async fn range_io_stats_count_coalesced_reads() {
+ let data = Bytes::from(vec![7u8; 100_000]);
+ let source: Arc<dyn FileRead> = TrackingRead::new(data.clone());
+ let mut reader = VindexFileReader::new(
+ source,
+ tokio::runtime::Handle::current(),
+ data.len() as u64,
+ "index".to_string(),
+ );
+ let stats = Arc::new(RangeIoStats::default());
+ reader.stats = Some(Arc::clone(&stats));
+
+ tokio::task::spawn_blocking(move || {
+ let mut first = [0u8; 4];
+ let mut second = [0u8; 4];
+ let mut third = [0u8; 4];
+ reader
+ .pread(&mut [
+ ReadRequest::new(0, &mut first),
+ ReadRequest::new(8, &mut second),
+ ReadRequest::new(20_000, &mut third),
+ ])
+ .unwrap();
+ })
+ .await
+ .unwrap();
+
+ let stats = stats.snapshot();
+ assert_eq!(stats.logical_ranges, 3);
+ assert_eq!(stats.requested_bytes, 12);
+ assert!(stats.file_read_calls < stats.logical_ranges);
+ assert!(stats.returned_bytes >= stats.requested_bytes);
+ assert_eq!(stats.read_ahead_hits, 0);
+ }
+
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn shared_permits_bound_reads_across_independent_readers() {
let data = Bytes::from(vec![8u8; 1024]);
@@ -600,7 +764,7 @@ mod tests {
active: AtomicUsize::new(0),
max_active: AtomicUsize::new(0),
});
- let permits = Arc::new(tokio::sync::Semaphore::new(1));
+ let permits = Arc::new(tokio::sync::Semaphore::new(0));
let make_reader = |path: &str| {
let source: Arc<dyn FileRead> = tracking.clone();
VindexFileReader::new_with_permits(
@@ -613,6 +777,9 @@ mod tests {
};
let mut first_reader = make_reader("first.index");
let mut second_reader = make_reader("second.index");
+ let stats = Arc::new(RangeIoStats::default());
+ first_reader.stats = Some(Arc::clone(&stats));
+ second_reader.stats = Some(Arc::clone(&stats));
assert!(Arc::ptr_eq(&first_reader.permits, &second_reader.permits));
let first = tokio::task::spawn_blocking(move || {
@@ -627,10 +794,15 @@ mod tests {
.pread(&mut [ReadRequest::new(128, &mut output)])
.unwrap();
});
+ tokio::time::sleep(Duration::from_millis(10)).await;
+ permits.add_permits(1);
first.await.unwrap();
second.await.unwrap();
assert_eq!(tracking.max_active.load(Ordering::SeqCst), 1);
+ let stats = stats.snapshot();
+ assert!(stats.range_permit_wait_nanos > 0);
+ assert!(stats.io_wait_nanos >= stats.range_permit_wait_nanos);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
diff --git a/crates/paimon/src/vindex/reader.rs
b/crates/paimon/src/vindex/reader.rs
index 55112dd8..f0cb0fe8 100644
--- a/crates/paimon/src/vindex/reader.rs
+++ b/crates/paimon/src/vindex/reader.rs
@@ -16,6 +16,7 @@
// under the License.
use crate::vector_search::{GlobalIndexIOMeta, VectorSearch};
+use crate::vindex::vector_search_timing_enabled;
use paimon_vindex_core::distance::MetricType;
use paimon_vindex_core::index::{
VectorIndexMetadata, VectorIndexReader as VIndexReader, VectorSearchParams,
@@ -25,6 +26,7 @@ use std::collections::BinaryHeap;
use std::collections::HashMap;
use std::io;
use std::sync::{Condvar, Mutex};
+use std::time::{Duration, Instant};
const DEFAULT_NPROBE: usize = 16;
const NPROBE_PARAMETER: &str = "ivf.nprobe";
@@ -91,6 +93,24 @@ fn native_batch_memory_reservation(index_parallelism: usize)
-> usize {
NATIVE_BATCH_PROCESS_WORKING_SET_BYTES / index_parallelism.max(1)
}
+#[derive(Clone, Copy, Default)]
+struct VindexLoadTiming {
+ vindex_open: Duration,
+ metadata: Duration,
+ optimize: Duration,
+}
+
+#[derive(Default)]
+struct VindexBatchStats {
+ native_chunk_queries: Vec<usize>,
+ scalar_chunk_count: usize,
+ max_chunk_size: usize,
+ memory_budget_bytes: usize,
+ batch_index_parallelism: usize,
+}
+
+type VindexBatchSearchResult = (Vec<Option<HashMap<u64, f32>>>,
Option<VindexBatchStats>);
+
trait ErasedSeekRead: Send {
fn pread_erased(&mut self, ranges: &mut [ReadRequest<'_>]) ->
io::Result<()>;
@@ -142,6 +162,9 @@ pub struct VindexVectorGlobalIndexReader {
batch_index_parallelism: usize,
reader: Option<VIndexReader<VindexInput>>,
metadata: Option<VectorIndexMetadata>,
+ timing_enabled: bool,
+ load_timing: VindexLoadTiming,
+ batch_stats: Option<VindexBatchStats>,
}
impl VindexVectorGlobalIndexReader {
@@ -152,6 +175,9 @@ impl VindexVectorGlobalIndexReader {
batch_index_parallelism: 1,
reader: None,
metadata: None,
+ timing_enabled: vector_search_timing_enabled(),
+ load_timing: VindexLoadTiming::default(),
+ batch_stats: None,
}
}
@@ -165,8 +191,10 @@ impl VindexVectorGlobalIndexReader {
vector_search: &VectorSearch,
stream_fn: impl FnOnce(&str) -> crate::Result<S>,
) -> crate::Result<Option<HashMap<u64, f32>>> {
- self.ensure_loaded(stream_fn, |_| Ok(()))?;
- self.search(vector_search)
+ Ok(self
+ .visit_batch_vector_search(std::slice::from_ref(vector_search),
stream_fn)?
+ .pop()
+ .expect("single vector search result"))
}
pub fn visit_batch_vector_search<S: SeekRead + 'static>(
@@ -174,28 +202,74 @@ impl VindexVectorGlobalIndexReader {
vector_searches: &[VectorSearch],
stream_fn: impl FnOnce(&str) -> crate::Result<S>,
) -> crate::Result<Vec<Option<HashMap<u64, f32>>>> {
- self.ensure_loaded(stream_fn, |_| Ok(()))?;
- self.search_batch(vector_searches)
- }
-
- #[cfg(test)]
- pub(crate) fn load<S: SeekRead + 'static>(
- &mut self,
- stream_fn: impl FnOnce(&str) -> crate::Result<S>,
- ) -> crate::Result<()> {
- self.ensure_loaded(stream_fn, |_| Ok(()))
+ self.visit_batch_vector_search_validated(vector_searches, stream_fn,
|_| Ok(()))
}
- pub(crate) fn load_validated<S, F>(
+ pub(crate) fn visit_batch_vector_search_validated<S, F>(
&mut self,
+ vector_searches: &[VectorSearch],
stream_fn: impl FnOnce(&str) -> crate::Result<S>,
validate: F,
- ) -> crate::Result<()>
+ ) -> crate::Result<Vec<Option<HashMap<u64, f32>>>>
where
S: SeekRead + 'static,
F: FnOnce(&VectorIndexMetadata) -> crate::Result<()>,
{
- self.ensure_loaded(stream_fn, validate)
+ let total_start = self.timing_enabled.then(Instant::now);
+ self.ensure_loaded(stream_fn, validate)?;
+ let search_start = self.timing_enabled.then(Instant::now);
+ let results = self.search_batch(vector_searches)?;
+ if let (Some(total_start), Some(search_start), Some(stats)) =
+ (total_start, search_start, self.batch_stats.as_ref())
+ {
+ let total = total_start.elapsed();
+ let native_search = search_start.elapsed();
+ let load = self
+ .load_timing
+ .vindex_open
+ .saturating_add(self.load_timing.metadata)
+ .saturating_add(self.load_timing.optimize);
+ let unattributed =
total.saturating_sub(load.saturating_add(native_search));
+ let chunk_queries = stats
+ .native_chunk_queries
+ .iter()
+ .map(usize::to_string)
+ .collect::<Vec<_>>()
+ .join(",");
+ let nprobe = self
+ .options
+ .get(NPROBE_PARAMETER)
+ .cloned()
+ .unwrap_or_else(|| DEFAULT_NPROBE.to_string());
+ log::debug!(
+ target: "paimon::vector_search",
+ "event=paimon_vindex_reader file={} nq={} nprobe={}
batch_index_parallelism={} memory_budget_bytes={} max_chunk_size={}
native_chunk_count={} native_chunk_queries={} scalar_chunk_count={}
total_ms={:.3} vindex_open_ms={:.3} metadata_ms={:.3} optimize_ms={:.3}
native_search_wall_ms={:.3} unattributed_ms={:.3}",
+ self.io_meta.file_path,
+ vector_searches.len(),
+ nprobe,
+ stats.batch_index_parallelism,
+ stats.memory_budget_bytes,
+ stats.max_chunk_size,
+ stats.native_chunk_queries.len(),
+ chunk_queries,
+ stats.scalar_chunk_count,
+ total.as_secs_f64() * 1000.0,
+ self.load_timing.vindex_open.as_secs_f64() * 1000.0,
+ self.load_timing.metadata.as_secs_f64() * 1000.0,
+ self.load_timing.optimize.as_secs_f64() * 1000.0,
+ native_search.as_secs_f64() * 1000.0,
+ unattributed.as_secs_f64() * 1000.0,
+ );
+ }
+ Ok(results)
+ }
+
+ #[cfg(test)]
+ pub(crate) fn load<S: SeekRead + 'static>(
+ &mut self,
+ stream_fn: impl FnOnce(&str) -> crate::Result<S>,
+ ) -> crate::Result<()> {
+ self.ensure_loaded(stream_fn, |_| Ok(()))
}
pub(crate) fn metadata(&self) -> crate::Result<&VectorIndexMetadata> {
@@ -207,7 +281,7 @@ impl VindexVectorGlobalIndexReader {
})
}
- pub(crate) fn search_batch(
+ fn search_batch(
&mut self,
vector_searches: &[VectorSearch],
) -> crate::Result<Vec<Option<HashMap<u64, f32>>>> {
@@ -225,15 +299,19 @@ impl VindexVectorGlobalIndexReader {
message: "vindex metadata not initialized".to_string(),
source: None,
})?;
- search_batch_vindex(
+ let (results, batch_stats) = search_batch_vindex(
reader,
metadata,
&self.options,
vector_searches,
self.batch_index_parallelism,
- )
+ self.timing_enabled,
+ )?;
+ self.batch_stats = batch_stats;
+ Ok(results)
}
+ #[cfg(test)]
fn search(&mut self, vector_search: &VectorSearch) ->
crate::Result<Option<HashMap<u64, f32>>> {
let reader = self
.reader
@@ -284,9 +362,11 @@ impl VindexVectorGlobalIndexReader {
O: FnOnce(&mut VIndexReader<VindexInput>) -> crate::Result<()>,
{
if self.reader.is_some() {
+ self.load_timing = VindexLoadTiming::default();
return validate(self.metadata()?);
}
+ let open_start = self.timing_enabled.then(Instant::now);
let source = stream_fn(&self.io_meta.file_path)?;
let mut reader =
VIndexReader::open(VindexInput::new(source)).map_err(|e| {
crate::Error::DataInvalid {
@@ -294,16 +374,27 @@ impl VindexVectorGlobalIndexReader {
source: Some(Box::new(e)),
}
})?;
+ let vindex_open = open_start.map_or(Duration::ZERO, |start|
start.elapsed());
+ let metadata_start = self.timing_enabled.then(Instant::now);
let metadata = reader.metadata();
+ let metadata_elapsed = metadata_start.map_or(Duration::ZERO, |start|
start.elapsed());
validate(&metadata)?;
+ let optimize_start = self.timing_enabled.then(Instant::now);
optimize(&mut reader)?;
+ let optimize_elapsed = optimize_start.map_or(Duration::ZERO, |start|
start.elapsed());
self.reader = Some(reader);
self.metadata = Some(metadata);
+ self.load_timing = VindexLoadTiming {
+ vindex_open,
+ metadata: metadata_elapsed,
+ optimize: optimize_elapsed,
+ };
Ok(())
}
}
+#[cfg(test)]
fn search_vindex(
reader: &mut VIndexReader<impl SeekRead>,
metadata: &VectorIndexMetadata,
@@ -406,10 +497,16 @@ fn search_batch_vindex(
options: &HashMap<String, String>,
vector_searches: &[VectorSearch],
index_parallelism: usize,
-) -> crate::Result<Vec<Option<HashMap<u64, f32>>>> {
+ timing_enabled: bool,
+) -> crate::Result<VindexBatchSearchResult> {
let mut results: Vec<Option<HashMap<u64, f32>>> =
(0..vector_searches.len()).map(|_| None).collect();
let mut groups: Vec<(PreparedSearch, Vec<usize>)> = Vec::new();
+ let mut batch_stats = timing_enabled.then(|| VindexBatchStats {
+ memory_budget_bytes:
native_batch_memory_reservation(index_parallelism),
+ batch_index_parallelism: index_parallelism,
+ ..VindexBatchStats::default()
+ });
for (index, search) in vector_searches.iter().enumerate() {
let Some(prepared) = prepare_search(metadata, options, search)? else {
@@ -424,8 +521,14 @@ fn search_batch_vindex(
for (prepared, indices) in groups {
let chunk_size = native_batch_chunk_size(metadata, &prepared,
index_parallelism);
+ if let Some(stats) = &mut batch_stats {
+ stats.max_chunk_size = stats.max_chunk_size.max(chunk_size);
+ }
for indices in indices.chunks(chunk_size) {
if indices.len() == 1 {
+ if let Some(stats) = &mut batch_stats {
+ stats.scalar_chunk_count += 1;
+ }
let index = indices[0];
let (labels, distances) =
execute_scalar_search(reader, &vector_searches[index],
&prepared)?;
@@ -435,6 +538,9 @@ fn search_batch_vindex(
}
continue;
}
+ if let Some(stats) = &mut batch_stats {
+ stats.native_chunk_queries.push(indices.len());
+ }
let reservation =
native_batch_chunk_working_set_bytes(metadata, &prepared,
indices.len());
@@ -486,7 +592,7 @@ fn search_batch_vindex(
}
}
- Ok(results)
+ Ok((results, batch_stats))
}
fn native_batch_chunk_size(