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(

Reply via email to