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, @@ -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> { + 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::>>()?; 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::(); + 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, } +#[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> { + 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::>(), - ) + 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::>>()?; + 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::(); + 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, index_search_limit: usize, ) -> crate::Result> { + 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> { +) -> crate::Result<(Vec, Option)> { + 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::>(); - 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> = 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 = 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, + stats: Option>, } 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> { + self.stats.clone() + } + fn validate_range(&self, pos: u64, len: usize) -> io::Result> { 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 = 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 = 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 = 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, + scalar_chunk_count: usize, + max_chunk_size: usize, + memory_budget_bytes: usize, + batch_index_parallelism: usize, +} + +type VindexBatchSearchResult = (Vec>>, Option); + 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>, metadata: Option, + timing_enabled: bool, + load_timing: VindexLoadTiming, + batch_stats: Option, } 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, ) -> crate::Result>> { - 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( @@ -174,28 +202,74 @@ impl VindexVectorGlobalIndexReader { vector_searches: &[VectorSearch], stream_fn: impl FnOnce(&str) -> crate::Result, ) -> crate::Result>>> { - self.ensure_loaded(stream_fn, |_| Ok(()))?; - self.search_batch(vector_searches) - } - - #[cfg(test)] - pub(crate) fn load( - &mut self, - stream_fn: impl FnOnce(&str) -> crate::Result, - ) -> crate::Result<()> { - self.ensure_loaded(stream_fn, |_| Ok(())) + self.visit_batch_vector_search_validated(vector_searches, stream_fn, |_| Ok(())) } - pub(crate) fn load_validated( + pub(crate) fn visit_batch_vector_search_validated( &mut self, + vector_searches: &[VectorSearch], stream_fn: impl FnOnce(&str) -> crate::Result, validate: F, - ) -> crate::Result<()> + ) -> crate::Result>>> 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::>() + .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( + &mut self, + stream_fn: impl FnOnce(&str) -> crate::Result, + ) -> 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>>> { @@ -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>> { let reader = self .reader @@ -284,9 +362,11 @@ impl VindexVectorGlobalIndexReader { O: FnOnce(&mut VIndexReader) -> 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, metadata: &VectorIndexMetadata, @@ -406,10 +497,16 @@ fn search_batch_vindex( options: &HashMap, vector_searches: &[VectorSearch], index_parallelism: usize, -) -> crate::Result>>> { + timing_enabled: bool, +) -> crate::Result { let mut results: Vec>> = (0..vector_searches.len()).map(|_| None).collect(); let mut groups: Vec<(PreparedSearch, Vec)> = 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(