From 356fbbc5b157ac08549c4bd8995a53b33ac7d1f8 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 18 Sep 2026 06:55:55 +0000 Subject: [PATCH 1/3] perf(index): overlap IVF query partition load and score Main's global-top-k path waits for 64 MiB / 128 partitions before the first score, and prepares only as many partitions as CPU cores. On a cold IVF_RQ query that serializes I/O and scoring. Gate four independent opts behind LANCE_IVF_QUERY_OPTS (default all): - first_wave: score after the first prepare wave - io_prepare: prepare concurrency follows I/O parallelism - parallel_files: load index.idx and auxiliary.idx together - overlap_prefilter: start partition I/O without waiting for prefilter Co-authored-by: Yang Cen --- rust/lance/src/index/vector/ivf/v2.rs | 274 +++++++++++++++++++++++--- 1 file changed, 246 insertions(+), 28 deletions(-) diff --git a/rust/lance/src/index/vector/ivf/v2.rs b/rust/lance/src/index/vector/ivf/v2.rs index d4f3134c2d4..0be26547b21 100644 --- a/rust/lance/src/index/vector/ivf/v2.rs +++ b/rust/lance/src/index/vector/ivf/v2.rs @@ -88,7 +88,7 @@ use prost::Message; use roaring::RoaringBitmap; use tokio::sync::mpsc; use tokio_stream::wrappers::ReceiverStream; -use tracing::{info, instrument}; +use tracing::{info, instrument, warn}; use uuid::Uuid; use super::{IvfIndexPartitionStatistics, IvfIndexStatistics, maybe_centroids_for_stats}; @@ -199,6 +199,101 @@ pub(crate) static GLOBAL_TOPK_CHUNK_BYTES: LazyLock = LazyLock::new(|| { /// expressible as a partition count: at most the prepare window plus two chunks. pub(crate) const GLOBAL_TOPK_CHUNK_MAX_PARTITIONS: usize = 128; +/// Query-time IVF load/score knobs that can be A/B'd independently. +/// +/// `LANCE_IVF_QUERY_OPTS` is a comma-separated list, `all`, or `none`. Unset +/// defaults to `all`. Known items: +/// - `first_wave`: score after the first prepare wave instead of waiting for +/// 64 MiB / 128 partitions. Later chunks keep the memory/dispatch bound. +/// - `io_prepare`: prepare concurrency follows I/O parallelism, not just CPU. +/// - `parallel_files`: load `index.idx` and `auxiliary.idx` concurrently. +/// - `overlap_prefilter`: start partition I/O without waiting for the prefilter. +const IVF_QUERY_OPTS_ENV: &str = "LANCE_IVF_QUERY_OPTS"; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct IvfQueryOpts { + first_wave: bool, + io_prepare: bool, + parallel_files: bool, + overlap_prefilter: bool, +} + +impl IvfQueryOpts { + const ALL: Self = Self { + first_wave: true, + io_prepare: true, + parallel_files: true, + overlap_prefilter: true, + }; + const NONE: Self = Self { + first_wave: false, + io_prepare: false, + parallel_files: false, + overlap_prefilter: false, + }; + + fn from_env() -> Self { + match std::env::var(IVF_QUERY_OPTS_ENV) { + Ok(raw) => Self::parse(&raw), + Err(_) => Self::ALL, + } + } + + fn parse(raw: &str) -> Self { + let raw = raw.trim(); + if raw.is_empty() || raw.eq_ignore_ascii_case("all") { + return Self::ALL; + } + if raw.eq_ignore_ascii_case("none") { + return Self::NONE; + } + let mut opts = Self::NONE; + for part in raw.split(',') { + match part.trim() { + "first_wave" => opts.first_wave = true, + "io_prepare" => opts.io_prepare = true, + "parallel_files" => opts.parallel_files = true, + "overlap_prefilter" => opts.overlap_prefilter = true, + "" => {} + other => { + warn!( + "ignoring unknown {IVF_QUERY_OPTS_ENV} item {other:?}; \ + expected first_wave, io_prepare, parallel_files, overlap_prefilter, all, or none" + ); + } + } + } + opts + } +} + +fn query_prepare_parallelism(io_parallelism: usize, opts: IvfQueryOpts) -> usize { + let cpu = get_num_compute_intensive_cpus().max(1); + if opts.io_prepare { + io_parallelism.max(1).max(cpu) + } else { + cpu + } +} + +fn scoring_chunk_limits( + opts: IvfQueryOpts, + is_first: bool, + prepare_parallelism: usize, +) -> (usize, usize) { + let chunk_bytes = *GLOBAL_TOPK_CHUNK_BYTES; + if is_first && opts.first_wave { + (chunk_bytes, prepare_parallelism.max(1)) + } else { + (chunk_bytes, GLOBAL_TOPK_CHUNK_MAX_PARTITIONS) + } +} + +#[cfg(test)] +fn global_topk_in_flight_bound(prepare_parallelism: usize) -> usize { + prepare_parallelism + 2 * GLOBAL_TOPK_CHUNK_MAX_PARTITIONS +} + /// Largest global-top-k heap converted to the result batch on the async task /// instead of a `spawn_cpu` dispatch; see the use site for the rationale. const GLOBAL_TOPK_INLINE_HEAP_LEN: usize = 4096; @@ -1320,19 +1415,20 @@ impl IVFIndex { } /// Pull prepared partitions off `prepared` until their pinned storage reaches - /// `chunk_bytes` or the chunk holds [`GLOBAL_TOPK_CHUNK_MAX_PARTITIONS`], - /// always taking at least one. `None` once the stream is exhausted; a failed - /// prepare ends the chunk (and the search) immediately. + /// `chunk_bytes` or the chunk holds `max_partitions`, always taking at least + /// one. `None` once the stream is exhausted; a failed prepare ends the chunk + /// (and the search) immediately. async fn next_scoring_chunk( prepared: &mut St, chunk_bytes: usize, + max_partitions: usize, ) -> Option>>> where St: Stream>> + Unpin, { let mut chunk = Vec::new(); let mut bytes = 0; - while bytes < chunk_bytes && chunk.len() < GLOBAL_TOPK_CHUNK_MAX_PARTITIONS { + while bytes < chunk_bytes && chunk.len() < max_partitions { match prepared.next().await { Some(Ok(prepared)) => { bytes += prepared.part_entry.size_bytes(); @@ -1740,11 +1836,11 @@ impl IVFIndex { } } - async fn load_partition_entry( + async fn load_partition_index( &self, partition_id: usize, io_stats: Option, - ) -> Result> { + ) -> Result { // `concat_batches` indexes the batches by this schema's field positions // without comparing the two, so the schema has to describe exactly what // was read: the full file schema over a projected read would index past @@ -1798,9 +1894,31 @@ impl IVFIndex { S::metadata_key().to_owned(), self.sub_index_metadata[partition_id].clone(), )?; - let idx = S::load(batch)?; - let storage = self.load_partition_storage(partition_id, io_stats).await?; - Ok(PartitionEntry::new(idx, storage)) + S::load(batch) + } + + async fn load_partition_entry( + &self, + partition_id: usize, + io_stats: Option, + ) -> Result> { + // IVF-PQ/HNSW keep sub-index rows in `index.idx` and codes in + // `auxiliary.idx`. Overlapping those two files hides the sequential + // wait; IVF-RQ's index file is typically empty, so the join is a + // no-op there. Disable with `LANCE_IVF_QUERY_OPTS` to A/B. + if IvfQueryOpts::from_env().parallel_files { + let (idx, storage) = tokio::try_join!( + self.load_partition_index(partition_id, io_stats.clone()), + self.load_partition_storage(partition_id, io_stats), + )?; + Ok(PartitionEntry::new(idx, storage)) + } else { + let idx = self + .load_partition_index(partition_id, io_stats.clone()) + .await?; + let storage = self.load_partition_storage(partition_id, io_stats).await?; + Ok(PartitionEntry::new(idx, storage)) + } } async fn materialize_prewarm_partition( @@ -2303,15 +2421,19 @@ impl VectorIndex for IVFInd ))); } - let prepare_parallelism = get_num_compute_intensive_cpus().max(1); + let opts = IvfQueryOpts::from_env(); + let prepare_parallelism = query_prepare_parallelism(self.io_parallelism, opts); let raw_query_context = self.prepare_rq_raw_query_context(&query.key)?; if control.is_none() && S::supports_global_topk_heap() { let heap_capacity = query.k * query.refine_factor.unwrap_or(1) as usize; - pre_filter.wait_for_ready().await?; + if !opts.overlap_prefilter { + pre_filter.wait_for_ready().await?; + } let prepare_index = self.clone(); let prepare_metrics = metrics.clone(); let prepare_raw_query_context = raw_query_context.clone(); + let overlap_prefilter = opts.overlap_prefilter; // Stream prepared partitions through scoring in chunks rather than // collecting all of them first. A prepared partition pins its whole // quantized storage, so collecting `nprobes` of them before scoring @@ -2332,20 +2454,33 @@ impl VectorIndex for IVFInd let metrics = prepare_metrics.clone(); let raw_query_context = prepare_raw_query_context.clone(); async move { - index - .prepare_partition_without_prefilter_wait( - part_id as usize, - &query, - pre_filter, - metrics.as_ref(), - raw_query_context, - ) - .await + if overlap_prefilter { + index + .prepare_partition( + part_id as usize, + &query, + pre_filter, + metrics.as_ref(), + raw_query_context, + ) + .await + } else { + index + .prepare_partition_without_prefilter_wait( + part_id as usize, + &query, + pre_filter, + metrics.as_ref(), + raw_query_context, + ) + .await + } } }) .buffered(prepare_parallelism) .fuse(); - let chunk_bytes = *GLOBAL_TOPK_CHUNK_BYTES; + let (first_bytes, first_max) = scoring_chunk_limits(opts, true, prepare_parallelism); + let (rest_bytes, rest_max) = scoring_chunk_limits(opts, false, prepare_parallelism); let use_query_residual = self.use_query_residual; let use_residual_scratch = self.use_residual_scratch; @@ -2356,7 +2491,14 @@ impl VectorIndex for IVFInd // the waiting (#7642). The heap is threaded through each dispatch so // scoring stays sequential, and a scored chunk is dropped before the // next one is scored. - let mut pending = Self::next_scoring_chunk(&mut prepared, chunk_bytes).await; + // + // The first chunk is capped at the prepare wave so scoring starts as + // soon as that wave is resident. Typical nprobes (20-80) of IVF_RQ + // partitions sit well under 64 MiB, so the old byte/128 cap waited + // for every probe before the first `spawn_cpu`. Later chunks keep + // the 64 MiB / 128 bound so a 4096-probe warm query does not pay a + // dispatch per wave. + let mut pending = Self::next_scoring_chunk(&mut prepared, first_bytes, first_max).await; while let Some(chunk) = pending { let chunk = chunk?; let search_metrics = metrics.clone(); @@ -2377,8 +2519,10 @@ impl VectorIndex for IVFInd })?; Ok(heap) }); - let (scored, next) = - futures::join!(score, Self::next_scoring_chunk(&mut prepared, chunk_bytes)); + let (scored, next) = futures::join!( + score, + Self::next_scoring_chunk(&mut prepared, rest_bytes, rest_max) + ); heap = scored?; pending = next; } @@ -2682,7 +2826,8 @@ impl VectorIndex for IVFInd // Streaming bounds resident partition storage to the load window plus one // chunk. `buffered` preserves the sorted load order above, so scoring order // (and thus the k-th-distance tie-break) stays deterministic. - let load_parallelism = get_num_compute_intensive_cpus().max(1); + let load_parallelism = + query_prepare_parallelism(self.io_parallelism, IvfQueryOpts::from_env()); let load_index = self.clone(); let load_metrics = metrics.clone(); let mut loaded_chunks = stream::iter(assignment_list) @@ -3034,6 +3179,76 @@ mod tests { const LIGHTWEIGHT_PQ_PARTITIONS: usize = 2; const LIGHTWEIGHT_PQ_SUB_VECTORS: usize = 4; + #[test] + fn test_ivf_query_opts_parse() { + assert_eq!(super::IvfQueryOpts::parse("all"), super::IvfQueryOpts::ALL); + assert_eq!(super::IvfQueryOpts::parse(""), super::IvfQueryOpts::ALL); + assert_eq!( + super::IvfQueryOpts::parse("none"), + super::IvfQueryOpts::NONE + ); + assert_eq!( + super::IvfQueryOpts::parse("first_wave"), + super::IvfQueryOpts { + first_wave: true, + ..super::IvfQueryOpts::NONE + } + ); + assert_eq!( + super::IvfQueryOpts::parse("io_prepare,parallel_files"), + super::IvfQueryOpts { + io_prepare: true, + parallel_files: true, + ..super::IvfQueryOpts::NONE + } + ); + assert_eq!( + super::IvfQueryOpts::parse("overlap_prefilter"), + super::IvfQueryOpts { + overlap_prefilter: true, + ..super::IvfQueryOpts::NONE + } + ); + } + + #[test] + fn test_scoring_chunk_limits_first_wave() { + let opts = super::IvfQueryOpts { + first_wave: true, + ..super::IvfQueryOpts::NONE + }; + let prepare_parallelism = 8; + let (first_bytes, first_max) = super::scoring_chunk_limits(opts, true, prepare_parallelism); + let (rest_bytes, rest_max) = super::scoring_chunk_limits(opts, false, prepare_parallelism); + assert_eq!(first_bytes, *super::GLOBAL_TOPK_CHUNK_BYTES); + assert_eq!(first_max, prepare_parallelism); + assert_eq!(rest_bytes, *super::GLOBAL_TOPK_CHUNK_BYTES); + assert_eq!(rest_max, super::GLOBAL_TOPK_CHUNK_MAX_PARTITIONS); + } + + #[test] + fn test_scoring_chunk_limits_matches_main_when_disabled() { + let opts = super::IvfQueryOpts::NONE; + let (first_bytes, first_max) = super::scoring_chunk_limits(opts, true, 8); + let (rest_bytes, rest_max) = super::scoring_chunk_limits(opts, false, 8); + assert_eq!(first_bytes, rest_bytes); + assert_eq!(first_max, rest_max); + assert_eq!(first_max, super::GLOBAL_TOPK_CHUNK_MAX_PARTITIONS); + } + + #[test] + fn test_query_prepare_parallelism() { + let cpu = get_num_compute_intensive_cpus().max(1); + assert_eq!( + super::query_prepare_parallelism(64, super::IvfQueryOpts::NONE), + cpu + ); + assert_eq!( + super::query_prepare_parallelism(64, super::IvfQueryOpts::ALL), + 64.max(cpu) + ); + } + lance_testing::define_stage_event_progress!(RecordingProgress, IndexBuildProgress, Result<()>); #[test] @@ -6771,8 +6986,11 @@ mod tests { // far smaller than the chunk byte budget, so the partition cap is what // ends a chunk. Size the index so that the old collect-everything // behavior would clearly exceed the bound. - let in_flight_bound = - get_num_compute_intensive_cpus().max(1) + 2 * super::GLOBAL_TOPK_CHUNK_MAX_PARTITIONS; + let prepare_parallelism = super::query_prepare_parallelism( + lance_io::object_store::DEFAULT_LOCAL_IO_PARALLELISM, + super::IvfQueryOpts::from_env(), + ); + let in_flight_bound = super::global_topk_in_flight_bound(prepare_parallelism); let num_partitions = 2 * in_flight_bound; let test_dir = TempStrDir::default(); From e3c266ecf4cc0a1bbe84ca8ae1c0971a9d893e30 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 18 Sep 2026 08:12:07 +0000 Subject: [PATCH 2/3] perf(index): preload IVF_RQ aux side columns at open Load every auxiliary column except RaBitQ codes when the index is opened (`LANCE_IVF_QUERY_OPTS=preload_aux`). Query-time partition loads then read only the code columns and merge the resident side buffers. The preloaded batch is stashed on IvfIndexState so in-process reconstruction does not re-read the side columns. Not enabled by `all`. Co-authored-by: Yang Cen --- rust/lance-index/src/vector/storage.rs | 284 +++++++++++++++++- .../src/index/vector/ivf/partition_serde.rs | 1 + rust/lance/src/index/vector/ivf/v2.rs | 132 +++++++- 3 files changed, 406 insertions(+), 11 deletions(-) diff --git a/rust/lance-index/src/vector/storage.rs b/rust/lance-index/src/vector/storage.rs index bbb478f4862..8d59bbe8906 100644 --- a/rust/lance-index/src/vector/storage.rs +++ b/rust/lance-index/src/vector/storage.rs @@ -3,9 +3,12 @@ //! Vector Storage, holding (quantized) vectors and providing distance calculation. +use crate::vector::bq::storage::{ + RABIT_BLOCKED_EX_CODE_COLUMN, RABIT_CODE_COLUMN, RABIT_EX_CODE_COLUMN, +}; use crate::vector::quantizer::QuantizerStorage; use arrow::compute::concat_batches; -use arrow_array::{ArrayRef, RecordBatch}; +use arrow_array::{Array, ArrayRef, RecordBatch}; use arrow_schema::SchemaRef; use futures::prelude::stream::TryStreamExt; use lance_arrow::RecordBatchExt; @@ -13,7 +16,8 @@ use lance_core::deepsize::DeepSizeOf; use lance_core::utils::tokio::spawn_cpu; use lance_core::{Error, ROW_ID, Result}; use lance_encoding::decoder::FilterExpression; -use lance_file::reader::FileReader; +use lance_file::reader::{FileReader, ReaderProjection}; +use lance_file::versions::reader_projection_from_column_names; use lance_io::ReadBatchParams; use lance_io::scheduler::IoStats; use lance_linalg::distance::DistanceType; @@ -23,8 +27,8 @@ use std::{ borrow::Cow, collections::BinaryHeap, mem::size_of, - ops::{Deref, DerefMut}, - sync::Arc, + ops::{Deref, DerefMut, Range}, + sync::{Arc, OnceLock}, }; use crossbeam_queue::ArrayQueue; @@ -56,6 +60,99 @@ where spawn_cpu(materialize).await } +/// Quantization code columns that stay on disk until a partition is queried. +/// +/// IVF_RQ side columns (`_rowid`, residual factors) are much smaller than these +/// and can be loaded when the index is opened. +pub fn is_rq_code_column(name: &str) -> bool { + matches!( + name, + RABIT_CODE_COLUMN | RABIT_EX_CODE_COLUMN | RABIT_BLOCKED_EX_CODE_COLUMN + ) +} + +fn schema_names_matching( + schema: &lance_core::datatypes::Schema, + pred: impl Fn(&str) -> bool, +) -> Vec<&str> { + schema + .fields + .iter() + .map(|field| field.name.as_str()) + .filter(|name| pred(name)) + .collect() +} + +fn record_batch_heap_size(batch: &RecordBatch) -> usize { + batch + .columns() + .iter() + .map(|column| column.get_array_memory_size()) + .sum() +} + +fn merge_code_and_side_columns( + codes: &RecordBatch, + side: &RecordBatch, + full_schema: SchemaRef, +) -> Result { + if codes.num_rows() != side.num_rows() { + return Err(Error::index(format!( + "preloaded aux side columns have {} rows but codes have {} rows", + side.num_rows(), + codes.num_rows() + ))); + } + let columns = full_schema + .fields() + .iter() + .map(|field| { + if let Some(column) = codes.column_by_name(field.name()) { + Ok(column.clone()) + } else if let Some(column) = side.column_by_name(field.name()) { + Ok(column.clone()) + } else { + Err(Error::index(format!( + "missing column {} when merging preloaded IVF aux side columns", + field.name() + ))) + } + }) + .collect::>>()?; + Ok(RecordBatch::try_new(full_schema, columns)?) +} + +async fn read_projected_range( + reader: &FileReader, + range: Range, + projection: ReaderProjection, + io_stats: Option<&IoStats>, +) -> Result { + let arrow_schema = Arc::new(arrow_schema::Schema::from(projection.schema.as_ref())); + if range.is_empty() { + return Ok(RecordBatch::new_empty(arrow_schema)); + } + let reader = match io_stats { + Some(io_stats) => Cow::Owned(reader.with_io_stats(io_stats.recorder())), + None => Cow::Borrowed(reader), + }; + let batches = reader + .read_stream_projected( + ReadBatchParams::Range(range), + u32::MAX, + 8, + projection, + FilterExpression::no_filter(), + ) + .await? + .try_collect::>() + .await?; + if batches.is_empty() { + return Ok(RecordBatch::new_empty(arrow_schema)); + } + Ok(concat_batches(&arrow_schema, batches.iter())?) +} + fn compact_prewarm_batches(batches: Vec) -> Result { let schema = batches .first() @@ -552,11 +649,21 @@ pub struct IvfQuantizationStorage { ivf: IvfModel, frag_reuse_index: Option>, + /// Aux columns except quantization codes, loaded at index open. + /// Query-time partition loads then read only the code columns. + preloaded_side: OnceLock>, + codes_projection: OnceLock, } impl DeepSizeOf for IvfQuantizationStorage { fn deep_size_of_children(&self, context: &mut lance_core::deepsize::Context) -> usize { - self.metadata.deep_size_of_children(context) + self.ivf.deep_size_of_children(context) + self.metadata.deep_size_of_children(context) + + self.ivf.deep_size_of_children(context) + + self + .preloaded_side + .get() + .map(|batch| record_batch_heap_size(batch)) + .unwrap_or_default() } } @@ -623,6 +730,8 @@ impl IvfQuantizationStorage { metadata, ivf, frag_reuse_index, + preloaded_side: OnceLock::new(), + codes_projection: OnceLock::new(), }) } @@ -654,6 +763,8 @@ impl IvfQuantizationStorage { metadata, ivf, frag_reuse_index, + preloaded_side: OnceLock::new(), + codes_projection: OnceLock::new(), } } @@ -695,6 +806,74 @@ impl IvfQuantizationStorage { self.ivf.num_partitions() } + /// Load every auxiliary column except IVF_RQ code columns into memory. + /// + /// No-op when the file has no RaBitQ code columns. Query-time + /// [`Self::load_partition`] then reads only the remaining code columns. + pub async fn preload_non_code_columns(&self) -> Result<()> { + if self.preloaded_side.get().is_some() { + return Ok(()); + } + let schema = self.reader.schema(); + let code_names = schema_names_matching(schema, is_rq_code_column); + if code_names.is_empty() { + return Ok(()); + } + let side_names = schema_names_matching(schema, |name| !is_rq_code_column(name)); + if side_names.is_empty() { + return Ok(()); + } + + let side_projection = reader_projection_from_column_names( + self.reader.metadata().version(), + schema, + &side_names, + )?; + let codes_projection = reader_projection_from_column_names( + self.reader.metadata().version(), + schema, + &code_names, + )?; + let num_rows = usize::try_from(self.reader.num_rows()).map_err(|_| { + Error::index(format!( + "aux file row count {} does not fit in usize", + self.reader.num_rows() + )) + })?; + let batch = read_projected_range(&self.reader, 0..num_rows, side_projection, None).await?; + let _ = self.preloaded_side.set(Arc::new(batch)); + let _ = self.codes_projection.set(codes_projection); + Ok(()) + } + + pub fn set_preloaded_non_code_columns(&self, batch: Arc) { + let _ = self.preloaded_side.set(batch); + } + + pub fn preloaded_non_code_columns(&self) -> Option> { + self.preloaded_side.get().cloned() + } + + fn codes_projection(&self) -> Result> { + if self.preloaded_side.get().is_none() { + return Ok(None); + } + if let Some(projection) = self.codes_projection.get() { + return Ok(Some(projection.clone())); + } + let code_names = schema_names_matching(self.reader.schema(), is_rq_code_column); + if code_names.is_empty() { + return Ok(None); + } + let projection = reader_projection_from_column_names( + self.reader.metadata().version(), + self.reader.schema(), + &code_names, + )?; + let _ = self.codes_projection.set(projection.clone()); + Ok(Some(projection)) + } + /// Load a partition's quantization storage, optionally measuring the exact /// I/O it performs into `io_stats`. /// @@ -702,6 +881,9 @@ impl IvfQuantizationStorage { /// scheduler also records into the sink (a cheap clone that shares all /// cached metadata, so no file is re-opened). When `None`, the normal /// uninstrumented reader is used. + /// + /// If [`Self::preload_non_code_columns`] has run, only code columns are + /// read from disk and merged with the resident side columns. pub async fn load_partition( &self, part_id: usize, @@ -712,6 +894,26 @@ impl IvfQuantizationStorage { let schema = self.reader.schema(); let arrow_schema = arrow_schema::Schema::from(schema.as_ref()); RecordBatch::new_empty(Arc::new(arrow_schema)) + } else if let (Some(side), Some(codes_projection)) = + (self.preloaded_side.get(), self.codes_projection()?) + { + if range.end > side.num_rows() { + return Err(Error::index(format!( + "partition {part_id} row range {}..{} exceeds preloaded aux side columns ({} rows)", + range.start, + range.end, + side.num_rows() + ))); + } + let codes = read_projected_range( + &self.reader, + range.clone(), + codes_projection, + io_stats.as_ref(), + ) + .await?; + let side = side.slice(range.start, range.end - range.start); + merge_code_and_side_columns(&codes, &side, self.schema())? } else { let reader = match &io_stats { Some(io_stats) => Cow::Owned(self.reader.with_io_stats(io_stats.recorder())), @@ -771,8 +973,8 @@ impl IvfQuantizationStorage { #[cfg(test)] mod tests { use super::{ - QueryScratchCapacity, QueryScratchPool, compact_prewarm_batches, - spawn_prewarm_materialization, + QueryScratchCapacity, QueryScratchPool, compact_prewarm_batches, is_rq_code_column, + merge_code_and_side_columns, spawn_prewarm_materialization, }; use arrow_array::{Array, ArrayRef, RecordBatch, UInt64Array}; use lance_core::deepsize::DeepSizeOf; @@ -924,4 +1126,72 @@ mod tests { assert_eq!(scratch.u32.capacity(), 3); }); } + + #[test] + fn test_is_rq_code_column() { + assert!(is_rq_code_column("_rabit_codes")); + assert!(is_rq_code_column("__ex_codes")); + assert!(is_rq_code_column("__blocked_ex_codes")); + assert!(!is_rq_code_column("_rowid")); + assert!(!is_rq_code_column("__add_factors")); + assert!(!is_rq_code_column("__pq_code")); + } + + #[test] + fn test_merge_code_and_side_columns() { + let codes = RecordBatch::try_from_iter([( + "_rabit_codes", + Arc::new(UInt64Array::from_iter_values(10..13)) as ArrayRef, + )]) + .unwrap(); + let side = RecordBatch::try_from_iter([ + ( + "_rowid", + Arc::new(UInt64Array::from_iter_values(0..3)) as ArrayRef, + ), + ( + "__add_factors", + Arc::new(UInt64Array::from_iter_values(100..103)) as ArrayRef, + ), + ]) + .unwrap(); + let full_schema = arrow_schema::Schema::new(vec![ + arrow_schema::Field::new("_rowid", arrow_schema::DataType::UInt64, true), + arrow_schema::Field::new("_rabit_codes", arrow_schema::DataType::UInt64, true), + arrow_schema::Field::new("__add_factors", arrow_schema::DataType::UInt64, true), + ]); + let merged = merge_code_and_side_columns(&codes, &side, Arc::new(full_schema)).unwrap(); + assert_eq!(merged.num_columns(), 3); + assert_eq!(merged.num_rows(), 3); + assert_eq!( + merged + .column_by_name("_rowid") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values(), + &[0, 1, 2] + ); + assert_eq!( + merged + .column_by_name("_rabit_codes") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values(), + &[10, 11, 12] + ); + assert_eq!( + merged + .column_by_name("__add_factors") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values(), + &[100, 101, 102] + ); + } } diff --git a/rust/lance/src/index/vector/ivf/partition_serde.rs b/rust/lance/src/index/vector/ivf/partition_serde.rs index a1869357980..352e29ddc77 100644 --- a/rust/lance/src/index/vector/ivf/partition_serde.rs +++ b/rust/lance/src/index/vector/ivf/partition_serde.rs @@ -1143,6 +1143,7 @@ mod tests { index_file_size: 1024, aux_file_size: 512, rq_search_cache: empty_rabit_search_cache_cell(), + preloaded_aux_side: None, }; let entry = IvfStateEntryBox(Arc::new(state)); diff --git a/rust/lance/src/index/vector/ivf/v2.rs b/rust/lance/src/index/vector/ivf/v2.rs index 0be26547b21..cdb63b0cc98 100644 --- a/rust/lance/src/index/vector/ivf/v2.rs +++ b/rust/lance/src/index/vector/ivf/v2.rs @@ -20,7 +20,7 @@ use crate::index::vector::{IndexFileVersion, builder::index_type_string}; use crate::index::{PreFilter, vector::VectorIndex}; use arrow::compute::concat_batches; use arrow_arith::numeric::sub; -use arrow_array::{ArrayRef, Float32Array, RecordBatch, UInt32Array, UInt64Array}; +use arrow_array::{Array, ArrayRef, Float32Array, RecordBatch, UInt32Array, UInt64Array}; use arrow_schema::DataType; use async_trait::async_trait; use datafusion::error::{DataFusionError, Result as DataFusionResult}; @@ -126,6 +126,9 @@ pub(crate) struct IvfIndexState { pub(crate) aux_file_size: u64, /// Runtime-only cache, intentionally excluded from the CacheCodec wire format. pub(crate) rq_search_cache: RabitSearchCacheCell, + /// Runtime-only IVF_RQ aux columns except codes. Dropped on disk-cache + /// deserialize; in-process reconstruct keeps the open-time buffers. + pub(crate) preloaded_aux_side: Option>, } /// Number of prepared partitions handed to a single `spawn_cpu` dispatch on the @@ -208,6 +211,8 @@ pub(crate) const GLOBAL_TOPK_CHUNK_MAX_PARTITIONS: usize = 128; /// - `io_prepare`: prepare concurrency follows I/O parallelism, not just CPU. /// - `parallel_files`: load `index.idx` and `auxiliary.idx` concurrently. /// - `overlap_prefilter`: start partition I/O without waiting for the prefilter. +/// - `preload_aux`: at index open, load every IVF_RQ aux column except codes. +/// Not part of `all`; enable it explicitly for A/B. const IVF_QUERY_OPTS_ENV: &str = "LANCE_IVF_QUERY_OPTS"; #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -216,6 +221,7 @@ struct IvfQueryOpts { io_prepare: bool, parallel_files: bool, overlap_prefilter: bool, + preload_aux: bool, } impl IvfQueryOpts { @@ -224,12 +230,14 @@ impl IvfQueryOpts { io_prepare: true, parallel_files: true, overlap_prefilter: true, + preload_aux: false, }; const NONE: Self = Self { first_wave: false, io_prepare: false, parallel_files: false, overlap_prefilter: false, + preload_aux: false, }; fn from_env() -> Self { @@ -254,11 +262,13 @@ impl IvfQueryOpts { "io_prepare" => opts.io_prepare = true, "parallel_files" => opts.parallel_files = true, "overlap_prefilter" => opts.overlap_prefilter = true, + "preload_aux" => opts.preload_aux = true, "" => {} other => { warn!( "ignoring unknown {IVF_QUERY_OPTS_ENV} item {other:?}; \ - expected first_wave, io_prepare, parallel_files, overlap_prefilter, all, or none" + expected first_wave, io_prepare, parallel_files, overlap_prefilter, \ + preload_aux, all, or none" ); } } @@ -722,6 +732,17 @@ impl DeepSizeOf for IvfIndexState { .and_then(|cache| cache.as_ref().and_then(|cache| cache.as_ref().cloned())) .map(|cache| cache.rotated_centroids.len() * std::mem::size_of::()) .unwrap_or_default() + + self + .preloaded_aux_side + .as_ref() + .map(|batch| { + batch + .columns() + .iter() + .map(|column| column.get_array_memory_size()) + .sum::() + }) + .unwrap_or_default() } } @@ -829,6 +850,7 @@ impl CacheCodecImpl for IvfStateEntryBox { index_file_size: header.index_file_size, aux_file_size: header.aux_file_size, rq_search_cache: empty_rabit_search_cache_cell(), + preloaded_aux_side: None, }))) } @@ -1687,6 +1709,9 @@ impl IVFIndex { .map(|index| Arc::new(CompactFragReuseIndexHandle(index)) as Arc); let storage = IvfQuantizationStorage::try_new_with_remapper(storage_reader, frag_reuse_index).await?; + if IvfQueryOpts::from_env().preload_aux { + storage.preload_non_code_columns().await?; + } // Cache file metadata so reconstructions from IvfIndexState can skip // footer reads. @@ -1788,6 +1813,11 @@ impl IVFIndex { &self.prepared_partitions } + #[cfg(test)] + pub(crate) fn storage(&self) -> &IvfQuantizationStorage { + &self.storage + } + #[instrument(level = "debug", skip(self, metrics))] pub async fn load_partition( &self, @@ -2105,6 +2135,7 @@ impl IVFIndex { index_file_size: self.reader.metadata().file_size(), aux_file_size: self.storage.reader().metadata().file_size(), rq_search_cache: rabit_search_cache_cell(self.rq_search_cache.clone()), + preloaded_aux_side: self.storage.preloaded_non_code_columns(), })) } } @@ -3053,6 +3084,9 @@ async fn reconstruct_typed( state.distance_type, frag_reuse_index, ); + if let Some(side) = &state.preloaded_aux_side { + storage.set_preloaded_non_code_columns(side.clone()); + } let rq_search_cache = IVFIndex::::rq_search_cache_from_state(state, &storage)?; let parsed_uuid = Uuid::parse_str(&state.uuid) @@ -3098,12 +3132,17 @@ mod tests { use lance_arrow::FixedSizeListArrayExt; use lance_index::vector::bq::{ RQBuildParams, RQRotationType, + builder::RabitQuantizer, ex_dot::{blocked_ex_code_bytes, padded_query_len}, - storage::{RABIT_BLOCKED_EX_CODE_COLUMN, RabitQuantizationMetadata, RabitQueryEstimator}, + storage::{ + RABIT_BLOCKED_EX_CODE_COLUMN, RABIT_CODE_COLUMN, RabitQuantizationMetadata, + RabitQueryEstimator, + }, transform::{EX_ADD_FACTORS_COLUMN, EX_SCALE_FACTORS_COLUMN}, }; use lance_index::vector::ivf::storage::IvfModel; use lance_index::vector::storage::VectorStore; + use lance_index::vector::storage::is_rq_code_column; use lance_index::vector::v3::subindex::IvfSubIndex; use crate::dataset::{InsertBuilder, UpdateBuilder, WriteMode, WriteParams}; @@ -3154,7 +3193,7 @@ mod tests { use lance_index::{INDEX_AUXILIARY_FILE_NAME, metrics::NoOpMetricsCollector}; use lance_io::{ object_store::{ObjectStore, ObjectStoreParams, StorageOptionsAccessor}, - scheduler::{ScanScheduler, SchedulerConfig}, + scheduler::{IoStats, ScanScheduler, SchedulerConfig}, utils::CachedFileSize, }; use lance_linalg::distance::{DistanceType, multivec_distance}; @@ -3209,6 +3248,20 @@ mod tests { ..super::IvfQueryOpts::NONE } ); + assert_eq!( + super::IvfQueryOpts::parse("preload_aux"), + super::IvfQueryOpts { + preload_aux: true, + ..super::IvfQueryOpts::NONE + } + ); + assert_eq!( + super::IvfQueryOpts::parse("all"), + super::IvfQueryOpts { + preload_aux: false, + ..super::IvfQueryOpts::ALL + } + ); } #[test] @@ -6458,6 +6511,77 @@ mod tests { assert_rq_rotation_type(&dataset, rotation_type).await; } + #[tokio::test] + async fn test_ivf_rq_preload_aux_reads_only_codes_at_query() { + let test_dir = TempStrDir::default(); + let test_uri = test_dir.as_str(); + let (mut dataset, vectors) = generate_test_dataset::(test_uri, 0.0..1.0).await; + + let ivf_params = IvfBuildParams::new(4); + let rq_params = RQBuildParams::with_rotation_type(1, RQRotationType::Fast); + let params = VectorIndexParams::with_ivf_rq_params(DistanceType::L2, ivf_params, rq_params); + dataset + .create_index(&["vector"], IndexType::Vector, None, ¶ms, true) + .await + .unwrap(); + + let indices = dataset.load_indices().await.unwrap(); + let index = dataset + .open_vector_index("vector", &indices[0].uuid, &NoOpMetricsCollector) + .await + .unwrap(); + let ivf = index + .as_any() + .downcast_ref::>() + .expect("IVF_RQ index"); + + let part_id = (0..ivf.storage().num_partitions()) + .find(|&part_id| ivf.storage().partition_size(part_id) > 0) + .expect("non-empty partition"); + + let full_stats = IoStats::new(); + ivf.load_partition_storage(part_id, Some(full_stats.clone())) + .await + .unwrap(); + let full = full_stats.snapshot(); + + ivf.storage().preload_non_code_columns().await.unwrap(); + let side = ivf + .storage() + .preloaded_non_code_columns() + .expect("side columns should be resident after preload"); + assert!(side.num_rows() > 0); + for field in side.schema().fields() { + assert!( + !is_rq_code_column(field.name()), + "preloaded column {} should not be a code column", + field.name() + ); + } + assert!(side.column_by_name(ROW_ID).is_some()); + assert!(side.column_by_name(RABIT_CODE_COLUMN).is_none()); + + let codes_stats = IoStats::new(); + ivf.load_partition_storage(part_id, Some(codes_stats.clone())) + .await + .unwrap(); + let codes = codes_stats.snapshot(); + assert!( + codes.iops < full.iops, + "preloaded query IOPS {} should be below full-schema IOPS {}", + codes.iops, + full.iops + ); + assert!( + codes.bytes_read < full.bytes_read, + "preloaded query bytes {} should be below full-schema bytes {}", + codes.bytes_read, + full.bytes_read + ); + + test_recall::(params, 4, 0.5, "vector", &dataset, vectors).await; + } + #[rstest] #[case(4, DistanceType::L2, 0.9)] #[case(4, DistanceType::Cosine, 0.9)] From a634b0e18804293b82f4b3cf72d8afeb43b41ca8 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 18 Sep 2026 08:29:38 +0000 Subject: [PATCH 3/3] Revert "perf(index): preload IVF_RQ aux side columns at open" This reverts commit e3c266ecf4cc0a1bbe84ca8ae1c0971a9d893e30. --- rust/lance-index/src/vector/storage.rs | 284 +----------------- .../src/index/vector/ivf/partition_serde.rs | 1 - rust/lance/src/index/vector/ivf/v2.rs | 132 +------- 3 files changed, 11 insertions(+), 406 deletions(-) diff --git a/rust/lance-index/src/vector/storage.rs b/rust/lance-index/src/vector/storage.rs index 8d59bbe8906..bbb478f4862 100644 --- a/rust/lance-index/src/vector/storage.rs +++ b/rust/lance-index/src/vector/storage.rs @@ -3,12 +3,9 @@ //! Vector Storage, holding (quantized) vectors and providing distance calculation. -use crate::vector::bq::storage::{ - RABIT_BLOCKED_EX_CODE_COLUMN, RABIT_CODE_COLUMN, RABIT_EX_CODE_COLUMN, -}; use crate::vector::quantizer::QuantizerStorage; use arrow::compute::concat_batches; -use arrow_array::{Array, ArrayRef, RecordBatch}; +use arrow_array::{ArrayRef, RecordBatch}; use arrow_schema::SchemaRef; use futures::prelude::stream::TryStreamExt; use lance_arrow::RecordBatchExt; @@ -16,8 +13,7 @@ use lance_core::deepsize::DeepSizeOf; use lance_core::utils::tokio::spawn_cpu; use lance_core::{Error, ROW_ID, Result}; use lance_encoding::decoder::FilterExpression; -use lance_file::reader::{FileReader, ReaderProjection}; -use lance_file::versions::reader_projection_from_column_names; +use lance_file::reader::FileReader; use lance_io::ReadBatchParams; use lance_io::scheduler::IoStats; use lance_linalg::distance::DistanceType; @@ -27,8 +23,8 @@ use std::{ borrow::Cow, collections::BinaryHeap, mem::size_of, - ops::{Deref, DerefMut, Range}, - sync::{Arc, OnceLock}, + ops::{Deref, DerefMut}, + sync::Arc, }; use crossbeam_queue::ArrayQueue; @@ -60,99 +56,6 @@ where spawn_cpu(materialize).await } -/// Quantization code columns that stay on disk until a partition is queried. -/// -/// IVF_RQ side columns (`_rowid`, residual factors) are much smaller than these -/// and can be loaded when the index is opened. -pub fn is_rq_code_column(name: &str) -> bool { - matches!( - name, - RABIT_CODE_COLUMN | RABIT_EX_CODE_COLUMN | RABIT_BLOCKED_EX_CODE_COLUMN - ) -} - -fn schema_names_matching( - schema: &lance_core::datatypes::Schema, - pred: impl Fn(&str) -> bool, -) -> Vec<&str> { - schema - .fields - .iter() - .map(|field| field.name.as_str()) - .filter(|name| pred(name)) - .collect() -} - -fn record_batch_heap_size(batch: &RecordBatch) -> usize { - batch - .columns() - .iter() - .map(|column| column.get_array_memory_size()) - .sum() -} - -fn merge_code_and_side_columns( - codes: &RecordBatch, - side: &RecordBatch, - full_schema: SchemaRef, -) -> Result { - if codes.num_rows() != side.num_rows() { - return Err(Error::index(format!( - "preloaded aux side columns have {} rows but codes have {} rows", - side.num_rows(), - codes.num_rows() - ))); - } - let columns = full_schema - .fields() - .iter() - .map(|field| { - if let Some(column) = codes.column_by_name(field.name()) { - Ok(column.clone()) - } else if let Some(column) = side.column_by_name(field.name()) { - Ok(column.clone()) - } else { - Err(Error::index(format!( - "missing column {} when merging preloaded IVF aux side columns", - field.name() - ))) - } - }) - .collect::>>()?; - Ok(RecordBatch::try_new(full_schema, columns)?) -} - -async fn read_projected_range( - reader: &FileReader, - range: Range, - projection: ReaderProjection, - io_stats: Option<&IoStats>, -) -> Result { - let arrow_schema = Arc::new(arrow_schema::Schema::from(projection.schema.as_ref())); - if range.is_empty() { - return Ok(RecordBatch::new_empty(arrow_schema)); - } - let reader = match io_stats { - Some(io_stats) => Cow::Owned(reader.with_io_stats(io_stats.recorder())), - None => Cow::Borrowed(reader), - }; - let batches = reader - .read_stream_projected( - ReadBatchParams::Range(range), - u32::MAX, - 8, - projection, - FilterExpression::no_filter(), - ) - .await? - .try_collect::>() - .await?; - if batches.is_empty() { - return Ok(RecordBatch::new_empty(arrow_schema)); - } - Ok(concat_batches(&arrow_schema, batches.iter())?) -} - fn compact_prewarm_batches(batches: Vec) -> Result { let schema = batches .first() @@ -649,21 +552,11 @@ pub struct IvfQuantizationStorage { ivf: IvfModel, frag_reuse_index: Option>, - /// Aux columns except quantization codes, loaded at index open. - /// Query-time partition loads then read only the code columns. - preloaded_side: OnceLock>, - codes_projection: OnceLock, } impl DeepSizeOf for IvfQuantizationStorage { fn deep_size_of_children(&self, context: &mut lance_core::deepsize::Context) -> usize { - self.metadata.deep_size_of_children(context) - + self.ivf.deep_size_of_children(context) - + self - .preloaded_side - .get() - .map(|batch| record_batch_heap_size(batch)) - .unwrap_or_default() + self.metadata.deep_size_of_children(context) + self.ivf.deep_size_of_children(context) } } @@ -730,8 +623,6 @@ impl IvfQuantizationStorage { metadata, ivf, frag_reuse_index, - preloaded_side: OnceLock::new(), - codes_projection: OnceLock::new(), }) } @@ -763,8 +654,6 @@ impl IvfQuantizationStorage { metadata, ivf, frag_reuse_index, - preloaded_side: OnceLock::new(), - codes_projection: OnceLock::new(), } } @@ -806,74 +695,6 @@ impl IvfQuantizationStorage { self.ivf.num_partitions() } - /// Load every auxiliary column except IVF_RQ code columns into memory. - /// - /// No-op when the file has no RaBitQ code columns. Query-time - /// [`Self::load_partition`] then reads only the remaining code columns. - pub async fn preload_non_code_columns(&self) -> Result<()> { - if self.preloaded_side.get().is_some() { - return Ok(()); - } - let schema = self.reader.schema(); - let code_names = schema_names_matching(schema, is_rq_code_column); - if code_names.is_empty() { - return Ok(()); - } - let side_names = schema_names_matching(schema, |name| !is_rq_code_column(name)); - if side_names.is_empty() { - return Ok(()); - } - - let side_projection = reader_projection_from_column_names( - self.reader.metadata().version(), - schema, - &side_names, - )?; - let codes_projection = reader_projection_from_column_names( - self.reader.metadata().version(), - schema, - &code_names, - )?; - let num_rows = usize::try_from(self.reader.num_rows()).map_err(|_| { - Error::index(format!( - "aux file row count {} does not fit in usize", - self.reader.num_rows() - )) - })?; - let batch = read_projected_range(&self.reader, 0..num_rows, side_projection, None).await?; - let _ = self.preloaded_side.set(Arc::new(batch)); - let _ = self.codes_projection.set(codes_projection); - Ok(()) - } - - pub fn set_preloaded_non_code_columns(&self, batch: Arc) { - let _ = self.preloaded_side.set(batch); - } - - pub fn preloaded_non_code_columns(&self) -> Option> { - self.preloaded_side.get().cloned() - } - - fn codes_projection(&self) -> Result> { - if self.preloaded_side.get().is_none() { - return Ok(None); - } - if let Some(projection) = self.codes_projection.get() { - return Ok(Some(projection.clone())); - } - let code_names = schema_names_matching(self.reader.schema(), is_rq_code_column); - if code_names.is_empty() { - return Ok(None); - } - let projection = reader_projection_from_column_names( - self.reader.metadata().version(), - self.reader.schema(), - &code_names, - )?; - let _ = self.codes_projection.set(projection.clone()); - Ok(Some(projection)) - } - /// Load a partition's quantization storage, optionally measuring the exact /// I/O it performs into `io_stats`. /// @@ -881,9 +702,6 @@ impl IvfQuantizationStorage { /// scheduler also records into the sink (a cheap clone that shares all /// cached metadata, so no file is re-opened). When `None`, the normal /// uninstrumented reader is used. - /// - /// If [`Self::preload_non_code_columns`] has run, only code columns are - /// read from disk and merged with the resident side columns. pub async fn load_partition( &self, part_id: usize, @@ -894,26 +712,6 @@ impl IvfQuantizationStorage { let schema = self.reader.schema(); let arrow_schema = arrow_schema::Schema::from(schema.as_ref()); RecordBatch::new_empty(Arc::new(arrow_schema)) - } else if let (Some(side), Some(codes_projection)) = - (self.preloaded_side.get(), self.codes_projection()?) - { - if range.end > side.num_rows() { - return Err(Error::index(format!( - "partition {part_id} row range {}..{} exceeds preloaded aux side columns ({} rows)", - range.start, - range.end, - side.num_rows() - ))); - } - let codes = read_projected_range( - &self.reader, - range.clone(), - codes_projection, - io_stats.as_ref(), - ) - .await?; - let side = side.slice(range.start, range.end - range.start); - merge_code_and_side_columns(&codes, &side, self.schema())? } else { let reader = match &io_stats { Some(io_stats) => Cow::Owned(self.reader.with_io_stats(io_stats.recorder())), @@ -973,8 +771,8 @@ impl IvfQuantizationStorage { #[cfg(test)] mod tests { use super::{ - QueryScratchCapacity, QueryScratchPool, compact_prewarm_batches, is_rq_code_column, - merge_code_and_side_columns, spawn_prewarm_materialization, + QueryScratchCapacity, QueryScratchPool, compact_prewarm_batches, + spawn_prewarm_materialization, }; use arrow_array::{Array, ArrayRef, RecordBatch, UInt64Array}; use lance_core::deepsize::DeepSizeOf; @@ -1126,72 +924,4 @@ mod tests { assert_eq!(scratch.u32.capacity(), 3); }); } - - #[test] - fn test_is_rq_code_column() { - assert!(is_rq_code_column("_rabit_codes")); - assert!(is_rq_code_column("__ex_codes")); - assert!(is_rq_code_column("__blocked_ex_codes")); - assert!(!is_rq_code_column("_rowid")); - assert!(!is_rq_code_column("__add_factors")); - assert!(!is_rq_code_column("__pq_code")); - } - - #[test] - fn test_merge_code_and_side_columns() { - let codes = RecordBatch::try_from_iter([( - "_rabit_codes", - Arc::new(UInt64Array::from_iter_values(10..13)) as ArrayRef, - )]) - .unwrap(); - let side = RecordBatch::try_from_iter([ - ( - "_rowid", - Arc::new(UInt64Array::from_iter_values(0..3)) as ArrayRef, - ), - ( - "__add_factors", - Arc::new(UInt64Array::from_iter_values(100..103)) as ArrayRef, - ), - ]) - .unwrap(); - let full_schema = arrow_schema::Schema::new(vec![ - arrow_schema::Field::new("_rowid", arrow_schema::DataType::UInt64, true), - arrow_schema::Field::new("_rabit_codes", arrow_schema::DataType::UInt64, true), - arrow_schema::Field::new("__add_factors", arrow_schema::DataType::UInt64, true), - ]); - let merged = merge_code_and_side_columns(&codes, &side, Arc::new(full_schema)).unwrap(); - assert_eq!(merged.num_columns(), 3); - assert_eq!(merged.num_rows(), 3); - assert_eq!( - merged - .column_by_name("_rowid") - .unwrap() - .as_any() - .downcast_ref::() - .unwrap() - .values(), - &[0, 1, 2] - ); - assert_eq!( - merged - .column_by_name("_rabit_codes") - .unwrap() - .as_any() - .downcast_ref::() - .unwrap() - .values(), - &[10, 11, 12] - ); - assert_eq!( - merged - .column_by_name("__add_factors") - .unwrap() - .as_any() - .downcast_ref::() - .unwrap() - .values(), - &[100, 101, 102] - ); - } } diff --git a/rust/lance/src/index/vector/ivf/partition_serde.rs b/rust/lance/src/index/vector/ivf/partition_serde.rs index 352e29ddc77..a1869357980 100644 --- a/rust/lance/src/index/vector/ivf/partition_serde.rs +++ b/rust/lance/src/index/vector/ivf/partition_serde.rs @@ -1143,7 +1143,6 @@ mod tests { index_file_size: 1024, aux_file_size: 512, rq_search_cache: empty_rabit_search_cache_cell(), - preloaded_aux_side: None, }; let entry = IvfStateEntryBox(Arc::new(state)); diff --git a/rust/lance/src/index/vector/ivf/v2.rs b/rust/lance/src/index/vector/ivf/v2.rs index cdb63b0cc98..0be26547b21 100644 --- a/rust/lance/src/index/vector/ivf/v2.rs +++ b/rust/lance/src/index/vector/ivf/v2.rs @@ -20,7 +20,7 @@ use crate::index::vector::{IndexFileVersion, builder::index_type_string}; use crate::index::{PreFilter, vector::VectorIndex}; use arrow::compute::concat_batches; use arrow_arith::numeric::sub; -use arrow_array::{Array, ArrayRef, Float32Array, RecordBatch, UInt32Array, UInt64Array}; +use arrow_array::{ArrayRef, Float32Array, RecordBatch, UInt32Array, UInt64Array}; use arrow_schema::DataType; use async_trait::async_trait; use datafusion::error::{DataFusionError, Result as DataFusionResult}; @@ -126,9 +126,6 @@ pub(crate) struct IvfIndexState { pub(crate) aux_file_size: u64, /// Runtime-only cache, intentionally excluded from the CacheCodec wire format. pub(crate) rq_search_cache: RabitSearchCacheCell, - /// Runtime-only IVF_RQ aux columns except codes. Dropped on disk-cache - /// deserialize; in-process reconstruct keeps the open-time buffers. - pub(crate) preloaded_aux_side: Option>, } /// Number of prepared partitions handed to a single `spawn_cpu` dispatch on the @@ -211,8 +208,6 @@ pub(crate) const GLOBAL_TOPK_CHUNK_MAX_PARTITIONS: usize = 128; /// - `io_prepare`: prepare concurrency follows I/O parallelism, not just CPU. /// - `parallel_files`: load `index.idx` and `auxiliary.idx` concurrently. /// - `overlap_prefilter`: start partition I/O without waiting for the prefilter. -/// - `preload_aux`: at index open, load every IVF_RQ aux column except codes. -/// Not part of `all`; enable it explicitly for A/B. const IVF_QUERY_OPTS_ENV: &str = "LANCE_IVF_QUERY_OPTS"; #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -221,7 +216,6 @@ struct IvfQueryOpts { io_prepare: bool, parallel_files: bool, overlap_prefilter: bool, - preload_aux: bool, } impl IvfQueryOpts { @@ -230,14 +224,12 @@ impl IvfQueryOpts { io_prepare: true, parallel_files: true, overlap_prefilter: true, - preload_aux: false, }; const NONE: Self = Self { first_wave: false, io_prepare: false, parallel_files: false, overlap_prefilter: false, - preload_aux: false, }; fn from_env() -> Self { @@ -262,13 +254,11 @@ impl IvfQueryOpts { "io_prepare" => opts.io_prepare = true, "parallel_files" => opts.parallel_files = true, "overlap_prefilter" => opts.overlap_prefilter = true, - "preload_aux" => opts.preload_aux = true, "" => {} other => { warn!( "ignoring unknown {IVF_QUERY_OPTS_ENV} item {other:?}; \ - expected first_wave, io_prepare, parallel_files, overlap_prefilter, \ - preload_aux, all, or none" + expected first_wave, io_prepare, parallel_files, overlap_prefilter, all, or none" ); } } @@ -732,17 +722,6 @@ impl DeepSizeOf for IvfIndexState { .and_then(|cache| cache.as_ref().and_then(|cache| cache.as_ref().cloned())) .map(|cache| cache.rotated_centroids.len() * std::mem::size_of::()) .unwrap_or_default() - + self - .preloaded_aux_side - .as_ref() - .map(|batch| { - batch - .columns() - .iter() - .map(|column| column.get_array_memory_size()) - .sum::() - }) - .unwrap_or_default() } } @@ -850,7 +829,6 @@ impl CacheCodecImpl for IvfStateEntryBox { index_file_size: header.index_file_size, aux_file_size: header.aux_file_size, rq_search_cache: empty_rabit_search_cache_cell(), - preloaded_aux_side: None, }))) } @@ -1709,9 +1687,6 @@ impl IVFIndex { .map(|index| Arc::new(CompactFragReuseIndexHandle(index)) as Arc); let storage = IvfQuantizationStorage::try_new_with_remapper(storage_reader, frag_reuse_index).await?; - if IvfQueryOpts::from_env().preload_aux { - storage.preload_non_code_columns().await?; - } // Cache file metadata so reconstructions from IvfIndexState can skip // footer reads. @@ -1813,11 +1788,6 @@ impl IVFIndex { &self.prepared_partitions } - #[cfg(test)] - pub(crate) fn storage(&self) -> &IvfQuantizationStorage { - &self.storage - } - #[instrument(level = "debug", skip(self, metrics))] pub async fn load_partition( &self, @@ -2135,7 +2105,6 @@ impl IVFIndex { index_file_size: self.reader.metadata().file_size(), aux_file_size: self.storage.reader().metadata().file_size(), rq_search_cache: rabit_search_cache_cell(self.rq_search_cache.clone()), - preloaded_aux_side: self.storage.preloaded_non_code_columns(), })) } } @@ -3084,9 +3053,6 @@ async fn reconstruct_typed( state.distance_type, frag_reuse_index, ); - if let Some(side) = &state.preloaded_aux_side { - storage.set_preloaded_non_code_columns(side.clone()); - } let rq_search_cache = IVFIndex::::rq_search_cache_from_state(state, &storage)?; let parsed_uuid = Uuid::parse_str(&state.uuid) @@ -3132,17 +3098,12 @@ mod tests { use lance_arrow::FixedSizeListArrayExt; use lance_index::vector::bq::{ RQBuildParams, RQRotationType, - builder::RabitQuantizer, ex_dot::{blocked_ex_code_bytes, padded_query_len}, - storage::{ - RABIT_BLOCKED_EX_CODE_COLUMN, RABIT_CODE_COLUMN, RabitQuantizationMetadata, - RabitQueryEstimator, - }, + storage::{RABIT_BLOCKED_EX_CODE_COLUMN, RabitQuantizationMetadata, RabitQueryEstimator}, transform::{EX_ADD_FACTORS_COLUMN, EX_SCALE_FACTORS_COLUMN}, }; use lance_index::vector::ivf::storage::IvfModel; use lance_index::vector::storage::VectorStore; - use lance_index::vector::storage::is_rq_code_column; use lance_index::vector::v3::subindex::IvfSubIndex; use crate::dataset::{InsertBuilder, UpdateBuilder, WriteMode, WriteParams}; @@ -3193,7 +3154,7 @@ mod tests { use lance_index::{INDEX_AUXILIARY_FILE_NAME, metrics::NoOpMetricsCollector}; use lance_io::{ object_store::{ObjectStore, ObjectStoreParams, StorageOptionsAccessor}, - scheduler::{IoStats, ScanScheduler, SchedulerConfig}, + scheduler::{ScanScheduler, SchedulerConfig}, utils::CachedFileSize, }; use lance_linalg::distance::{DistanceType, multivec_distance}; @@ -3248,20 +3209,6 @@ mod tests { ..super::IvfQueryOpts::NONE } ); - assert_eq!( - super::IvfQueryOpts::parse("preload_aux"), - super::IvfQueryOpts { - preload_aux: true, - ..super::IvfQueryOpts::NONE - } - ); - assert_eq!( - super::IvfQueryOpts::parse("all"), - super::IvfQueryOpts { - preload_aux: false, - ..super::IvfQueryOpts::ALL - } - ); } #[test] @@ -6511,77 +6458,6 @@ mod tests { assert_rq_rotation_type(&dataset, rotation_type).await; } - #[tokio::test] - async fn test_ivf_rq_preload_aux_reads_only_codes_at_query() { - let test_dir = TempStrDir::default(); - let test_uri = test_dir.as_str(); - let (mut dataset, vectors) = generate_test_dataset::(test_uri, 0.0..1.0).await; - - let ivf_params = IvfBuildParams::new(4); - let rq_params = RQBuildParams::with_rotation_type(1, RQRotationType::Fast); - let params = VectorIndexParams::with_ivf_rq_params(DistanceType::L2, ivf_params, rq_params); - dataset - .create_index(&["vector"], IndexType::Vector, None, ¶ms, true) - .await - .unwrap(); - - let indices = dataset.load_indices().await.unwrap(); - let index = dataset - .open_vector_index("vector", &indices[0].uuid, &NoOpMetricsCollector) - .await - .unwrap(); - let ivf = index - .as_any() - .downcast_ref::>() - .expect("IVF_RQ index"); - - let part_id = (0..ivf.storage().num_partitions()) - .find(|&part_id| ivf.storage().partition_size(part_id) > 0) - .expect("non-empty partition"); - - let full_stats = IoStats::new(); - ivf.load_partition_storage(part_id, Some(full_stats.clone())) - .await - .unwrap(); - let full = full_stats.snapshot(); - - ivf.storage().preload_non_code_columns().await.unwrap(); - let side = ivf - .storage() - .preloaded_non_code_columns() - .expect("side columns should be resident after preload"); - assert!(side.num_rows() > 0); - for field in side.schema().fields() { - assert!( - !is_rq_code_column(field.name()), - "preloaded column {} should not be a code column", - field.name() - ); - } - assert!(side.column_by_name(ROW_ID).is_some()); - assert!(side.column_by_name(RABIT_CODE_COLUMN).is_none()); - - let codes_stats = IoStats::new(); - ivf.load_partition_storage(part_id, Some(codes_stats.clone())) - .await - .unwrap(); - let codes = codes_stats.snapshot(); - assert!( - codes.iops < full.iops, - "preloaded query IOPS {} should be below full-schema IOPS {}", - codes.iops, - full.iops - ); - assert!( - codes.bytes_read < full.bytes_read, - "preloaded query bytes {} should be below full-schema bytes {}", - codes.bytes_read, - full.bytes_read - ); - - test_recall::(params, 4, 0.5, "vector", &dataset, vectors).await; - } - #[rstest] #[case(4, DistanceType::L2, 0.9)] #[case(4, DistanceType::Cosine, 0.9)]