diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 1d45d49a..1f8b4705 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -28,6 +28,7 @@ const VECTOR_INDEX_SEARCH_MODE_OPTION: &str = "vector-index.search-mode"; const FULL_TEXT_INDEX_SEARCH_MODE_OPTION: &str = "full-text-index.search-mode"; const GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION: &str = "global-index.row-count-per-shard"; const GLOBAL_INDEX_THREAD_NUM_OPTION: &str = "global-index.thread-num"; +const GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION: &str = "global-index.vindex.read-thread-num"; const GLOBAL_INDEX_COLUMN_UPDATE_ACTION_OPTION: &str = "global-index.column-update-action"; const SORTED_INDEX_RECORDS_PER_RANGE_OPTION: &str = "sorted-index.records-per-range"; const BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE_OPTION: &str = "btree-index.fallback-scan-max-size"; @@ -133,6 +134,8 @@ const DYNAMIC_BUCKET_TARGET_ROW_NUM_OPTION: &str = "dynamic-bucket.target-row-nu const DEFAULT_DYNAMIC_BUCKET_TARGET_ROW_NUM: i64 = 200_000; const DEFAULT_GLOBAL_INDEX_ROW_COUNT_PER_SHARD: i64 = 100_000; const DEFAULT_GLOBAL_INDEX_THREAD_NUM: i64 = 32; +pub(crate) const DEFAULT_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM: usize = 64; +const MAX_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM: i64 = tokio::sync::Semaphore::MAX_PERMITS as i64; const MAX_GLOBAL_INDEX_THREAD_NUM: i64 = { let tokio_max = (usize::MAX >> 3) as u64; let i32_max = i32::MAX as u64; @@ -696,13 +699,15 @@ impl<'a> CoreOptions<'a> { Ok(value) } - /// Maximum number of concurrent tasks for global-index I/O, mirroring Java + /// Maximum number of concurrent global-index search tasks, mirroring Java /// `CoreOptions.GLOBAL_INDEX_THREAD_NUM` (key `global-index.thread-num`, /// default 32). Used as the per-operation fan-out limit for sorted BTree and /// bitmap shard reads, global-index vector search, and primary-key vector - /// search. A value of `1` reproduces strict sequential execution. A - /// non-positive value, or one above [`MAX_GLOBAL_INDEX_THREAD_NUM`], is a - /// misconfiguration and fails loud rather than being silently clamped. + /// search. Vindex file range reads use + /// [`Self::global_index_vindex_read_thread_num`] instead. A value of `1` + /// makes these search tasks sequential, but does not serialize Vindex range + /// reads. A non-positive value, or one above [`MAX_GLOBAL_INDEX_THREAD_NUM`], + /// is a misconfiguration and fails loud rather than being silently clamped. pub fn global_index_thread_num(&self) -> crate::Result { let value = self .parse_i64_option(GLOBAL_INDEX_THREAD_NUM_OPTION)? @@ -728,6 +733,36 @@ impl<'a> CoreOptions<'a> { Ok(value as usize) } + /// Maximum number of concurrent range reads shared by Vindex readers in one + /// search operation (key `global-index.vindex.read-thread-num`, default 64). + /// This is independent of [`Self::global_index_thread_num`]. + pub fn global_index_vindex_read_thread_num(&self) -> crate::Result { + let value = self + .parse_i64_option(GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION)? + .unwrap_or(DEFAULT_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM as i64); + if value <= 0 { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{}' must be greater than 0, got: {}", + GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION, value + ), + source: None, + }); + } + if value > MAX_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{}' must not exceed {}, got: {}", + GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION, + MAX_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM, + value + ), + source: None, + }); + } + Ok(value as usize) + } + pub fn sorted_index_records_per_range(&self) -> crate::Result { let value = self .parse_i64_option(SORTED_INDEX_RECORDS_PER_RANGE_OPTION)? @@ -1520,6 +1555,10 @@ mod tests { 100_000 ); assert_eq!(core_options.global_index_thread_num().unwrap(), 32); + assert_eq!( + core_options.global_index_vindex_read_thread_num().unwrap(), + 64 + ); assert_eq!( core_options.sorted_index_records_per_range().unwrap(), 100_000 @@ -1764,6 +1803,48 @@ mod tests { ); } + #[test] + fn test_global_index_vindex_read_thread_num_default_and_custom() { + assert_eq!( + CoreOptions::new(&HashMap::new()) + .global_index_vindex_read_thread_num() + .unwrap(), + 64 + ); + + for value in [32, 64] { + let options = HashMap::from([( + GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION.to_string(), + value.to_string(), + )]); + assert_eq!( + CoreOptions::new(&options) + .global_index_vindex_read_thread_num() + .unwrap(), + value + ); + } + } + + #[test] + fn test_global_index_vindex_read_thread_num_rejects_invalid_values() { + for value in [ + "0".to_string(), + "abc".to_string(), + (MAX_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM + 1).to_string(), + ] { + let options = HashMap::from([( + GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION.to_string(), + value, + )]); + let err = CoreOptions::new(&options) + .global_index_vindex_read_thread_num() + .expect_err("invalid vindex.read-thread-num should fail"); + assert!(matches!(err, crate::Error::DataInvalid { message, .. } + if message.contains(GLOBAL_INDEX_VINDEX_READ_THREAD_NUM_OPTION))); + } + } + #[test] fn test_sorted_index_records_per_range_rejects_invalid_values() { for value in ["0", "-1", "abc"] { diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 5919b4c7..06f34781 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -58,7 +58,7 @@ 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::{RangeIoStats, VindexFileReader}; +use crate::vindex::range_reader::{RangeIoStats, RangeReadLimiter, VindexFileReader}; use crate::vindex::reader::VindexVectorGlobalIndexReader; use crate::vindex::{is_vindex_index_type, vector_search_timing_enabled, VindexVectorIndexOptions}; use arrow_array::{Array, FixedSizeListArray, Float32Array, Int64Array, ListArray, RecordBatch}; @@ -144,7 +144,7 @@ fn log_vindex_range_io_stats(file: &str, query_count: usize, stats: &RangeIoStat 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}", + "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} peak_in_flight_reads={} read_many_merged_ranges={} read_many_chunks={} read_many_chunk_size_sum={} read_many_chunk_size_min={} read_many_chunk_size_max={}", file, query_count, stats.logical_ranges, @@ -154,9 +154,26 @@ fn log_vindex_range_io_stats(file: &str, query_count: usize, stats: &RangeIoStat 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, + stats.peak_in_flight_reads, + stats.read_many_merged_ranges, + stats.read_many_chunks, + stats.read_many_chunk_size_sum, + stats.read_many_chunk_size_min, + stats.read_many_chunk_size_max, ); } +fn vindex_concurrency_limits( + core_options: &CoreOptions<'_>, + entry_count: usize, + max_concurrency: usize, +) -> crate::Result<(usize, usize)> { + Ok(( + vindex_index_parallelism(entry_count, max_concurrency), + core_options.global_index_vindex_read_thread_num()?, + )) +} + pub struct VectorSearchBuilder<'a> { table: &'a Table, vector_column: Option, @@ -844,15 +861,16 @@ async fn plan_and_search_pk_candidates_batch( source: None, } })?; - let batch_index_parallelism = match backend { - VectorIndexBackend::Vindex => vindex_index_parallelism( + let (batch_index_parallelism, range_read_concurrency) = match backend { + VectorIndexBackend::Vindex => vindex_concurrency_limits( + core, plan.splits .iter() .map(|split| split.ann_segments.len()) .sum(), concurrency, - ), - VectorIndexBackend::Lumina => 1, + )?, + VectorIndexBackend::Lumina => (1, 0), }; // Production data-file reader, mirroring `table_read.rs::new_data_file_reader` @@ -878,11 +896,14 @@ async fn plan_and_search_pk_candidates_batch( let field_name = pk_col.to_string(); let loader_io = table.file_io().clone(); - let loader_range_read_permits = Arc::new(tokio::sync::Semaphore::new(concurrency)); + let loader_range_read_limiter = match backend { + VectorIndexBackend::Vindex => Some(RangeReadLimiter::new(range_read_concurrency)), + VectorIndexBackend::Lumina => None, + }; let loader: crate::vindex::pkvector::ann::SourceSegmentLoader = Box::new( move |segment: &BucketAnnSegment| { let io = loader_io.clone(); - let range_read_permits = Arc::clone(&loader_range_read_permits); + let range_read_limiter = loader_range_read_limiter.clone(); let path = segment.path.clone(); let file_size = segment.file_size; Box::pin(async move { @@ -908,10 +929,10 @@ async fn plan_and_search_pk_candidates_batch( source: None, })?; Ok(AnnSegmentSource::Vindex( - VindexFileReader::new_with_permits( + VindexFileReader::new_with_limiter( Arc::new(file_reader), current_tokio_runtime_handle()?, - range_read_permits, + range_read_limiter.expect("Vindex range-read limiter"), file_size, path, ), @@ -1629,18 +1650,24 @@ async fn evaluate_batch_vector_search( }); } ensure_global_index_executor_capacity(concurrency); - let range_read_permits = Arc::new(tokio::sync::Semaphore::new(concurrency)); - let batch_index_parallelism = vindex_index_parallelism( - vector_entries - .iter() - .filter(|entry| is_vindex_index_type(&entry.index_file.index_type)) - .count(), - concurrency, - ); + let vindex_entry_count = vector_entries + .iter() + .filter(|entry| is_vindex_index_type(&entry.index_file.index_type)) + .count(); + let (batch_index_parallelism, range_read_limiter) = if vindex_entry_count == 0 { + (1, None) + } else { + let (index_parallelism, range_read_concurrency) = + vindex_concurrency_limits(&core_options, vindex_entry_count, concurrency)?; + ( + index_parallelism, + Some(RangeReadLimiter::new(range_read_concurrency)), + ) + }; let futures: Vec<_> = vector_entries .into_iter() .map(|entry| { - let range_read_permits = Arc::clone(&range_read_permits); + let range_read_limiter = range_read_limiter.clone(); let global_meta = entry.index_file.global_index_meta.as_ref().unwrap(); let backend = VectorIndexBackend::from_index_type(&entry.index_file.index_type) .expect("filtered vector index type"); @@ -1721,10 +1748,10 @@ async fn evaluate_batch_vector_search( })?; file_reader_open = file_reader_open_start .map_or(Duration::ZERO, |start| start.elapsed()); - let source = VindexFileReader::new_with_permits( + let source = VindexFileReader::new_with_limiter( Arc::new(file_reader), runtime, - range_read_permits, + range_read_limiter.expect("Vindex range-read limiter"), file_size, file_name.clone(), ); @@ -3446,11 +3473,25 @@ mod tests { } #[test] - fn vindex_batch_parallelism_tracks_active_entries() { - assert_eq!(vindex_index_parallelism(1, 1), 1); - assert_eq!(vindex_index_parallelism(1, 64), 1); - assert_eq!(vindex_index_parallelism(8, 4), 4); - assert_eq!(vindex_index_parallelism(4, 8), 4); + fn vindex_concurrency_limits_are_independent() { + let default_options = HashMap::new(); + let default_core = CoreOptions::new(&default_options); + assert_eq!( + vindex_concurrency_limits(&default_core, 1, 32).unwrap(), + (1, 64) + ); + assert_eq!( + vindex_concurrency_limits(&default_core, 8, 4).unwrap(), + (4, 64) + ); + + let options = HashMap::from([( + "global-index.vindex.read-thread-num".to_string(), + "48".to_string(), + )]); + let core = CoreOptions::new(&options); + assert_eq!(vindex_concurrency_limits(&core, 1, 32).unwrap(), (1, 48)); + assert_eq!(vindex_concurrency_limits(&core, 8, 4).unwrap(), (4, 48)); } fn make_field(id: i32, name: &str) -> DataField { diff --git a/crates/paimon/src/vindex/executor.rs b/crates/paimon/src/vindex/executor.rs index 173b5d7d..b604ae61 100644 --- a/crates/paimon/src/vindex/executor.rs +++ b/crates/paimon/src/vindex/executor.rs @@ -240,7 +240,7 @@ fn max_physical_worker_count() -> usize { /// Shared global-index executor for synchronous search and range-I/O waits. It /// starts at the machine's available parallelism and grows lazily, but keeps the /// configured logical fan-out separate from a physical cap of four workers per -/// CPU (and at least the default 32 I/O workers). Growth workers expire after the +/// CPU (and at least 32 I/O workers). Growth workers expire after the /// same one-minute idle interval used by Java's `GlobalIndexReadThreadPool`. /// Smaller query limits are enforced by each query's bounded job scheduler. fn global_executor() -> &'static GlobalIndexExecutor { diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 0f5d965c..d5235369 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -16,20 +16,21 @@ // under the License. use crate::io::FileRead; +#[cfg(test)] +use crate::spec::DEFAULT_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM; use crate::vindex::vector_search_timing_enabled; use bytes::Bytes; -use futures::future::try_join_all; +use futures::{stream, StreamExt}; 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::sync::Arc; use std::time::Instant; const SCALAR_READ_MAX: usize = 64; const SCALAR_READ_AHEAD: u64 = 64 * 1024; const RANGE_COALESCE_GAP: u64 = 16 * 1024; -const RANGE_READ_CONCURRENCY: usize = 32; struct CachedRange { start: u64, @@ -57,6 +58,14 @@ struct MergedRange { requested_bytes: u64, } +struct InFlightRead<'a>(&'a AtomicU64); + +impl Drop for InFlightRead<'_> { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::Relaxed); + } +} + #[derive(Debug, Default)] pub(crate) struct RangeIoStats { logical_ranges: AtomicU64, @@ -66,9 +75,16 @@ pub(crate) struct RangeIoStats { read_ahead_hits: AtomicU64, io_wait_nanos: AtomicU64, range_permit_wait_nanos: AtomicU64, + in_flight_reads: AtomicU64, + peak_in_flight_reads: AtomicU64, + read_many_merged_ranges: AtomicU64, + read_many_chunks: AtomicU64, + read_many_chunk_size_sum: AtomicU64, + read_many_chunk_size_min: AtomicU64, + read_many_chunk_size_max: AtomicU64, } -#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +#[derive(Clone, Debug, Default, PartialEq, Eq)] pub(crate) struct RangeIoStatsSnapshot { pub(crate) logical_ranges: u64, pub(crate) requested_bytes: u64, @@ -77,6 +93,12 @@ pub(crate) struct RangeIoStatsSnapshot { pub(crate) read_ahead_hits: u64, pub(crate) io_wait_nanos: u64, pub(crate) range_permit_wait_nanos: u64, + pub(crate) peak_in_flight_reads: u64, + pub(crate) read_many_merged_ranges: u64, + pub(crate) read_many_chunks: u64, + pub(crate) read_many_chunk_size_sum: u64, + pub(crate) read_many_chunk_size_min: u64, + pub(crate) read_many_chunk_size_max: u64, } impl RangeIoStats { @@ -89,17 +111,50 @@ impl RangeIoStats { 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), + peak_in_flight_reads: self.peak_in_flight_reads.load(Ordering::Relaxed), + read_many_merged_ranges: self.read_many_merged_ranges.load(Ordering::Relaxed), + read_many_chunks: self.read_many_chunks.load(Ordering::Relaxed), + read_many_chunk_size_sum: self.read_many_chunk_size_sum.load(Ordering::Relaxed), + read_many_chunk_size_min: self.read_many_chunk_size_min.load(Ordering::Relaxed), + read_many_chunk_size_max: self.read_many_chunk_size_max.load(Ordering::Relaxed), + } + } +} + +#[derive(Clone)] +pub(crate) struct RangeReadLimiter { + io_permits: Arc, + response_permits: Arc, + io_limit: usize, + response_limit: usize, +} + +impl RangeReadLimiter { + pub(crate) fn new(io_limit: usize) -> Self { + let response_limit = io_limit + .saturating_mul(2) + .min(tokio::sync::Semaphore::MAX_PERMITS); + Self { + io_permits: Arc::new(tokio::sync::Semaphore::new(io_limit)), + response_permits: Arc::new(tokio::sync::Semaphore::new(response_limit)), + io_limit, + response_limit, } } } +struct RangeResponse { + data: Bytes, + _permit: tokio::sync::OwnedSemaphorePermit, +} + /// 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. pub(crate) struct VindexFileReader { reader: Arc, runtime: tokio::runtime::Handle, - permits: Arc, + limiter: RangeReadLimiter, file_size: u64, path: String, scalar_cache: Option, @@ -114,26 +169,26 @@ impl VindexFileReader { file_size: u64, path: String, ) -> Self { - Self::new_with_permits( + Self::new_with_limiter( reader, runtime, - Arc::new(tokio::sync::Semaphore::new(RANGE_READ_CONCURRENCY)), + RangeReadLimiter::new(DEFAULT_GLOBAL_INDEX_VINDEX_READ_THREAD_NUM), file_size, path, ) } - pub(crate) fn new_with_permits( + pub(crate) fn new_with_limiter( reader: Arc, runtime: tokio::runtime::Handle, - permits: Arc, + limiter: RangeReadLimiter, file_size: u64, path: String, ) -> Self { Self { reader, runtime, - permits, + limiter, file_size, path, scalar_cache: None, @@ -186,56 +241,81 @@ impl VindexFileReader { .saturating_add(SCALAR_READ_AHEAD) .max(range.end) .min(self.file_size); - let data = self.fetch_exact(range.start..read_end)?; - buf.copy_from_slice(&data[..buf.len()]); + let response = self.fetch_exact(range.start..read_end)?; + buf.copy_from_slice(&response.data[..buf.len()]); self.scalar_cache = Some(CachedRange { start: range.start, - data, + data: response.data, }); return Ok(()); } - let data = self.fetch_exact(range)?; - buf.copy_from_slice(&data); + let response = self.fetch_exact(range)?; + buf.copy_from_slice(&response.data); Ok(()) } - fn fetch_exact(&self, range: Range) -> io::Result { - let mut results = self.fetch_range_batch(std::slice::from_ref(&range))?; - Ok(results.pop().expect("one requested range")) + fn fetch_exact(&self, range: Range) -> io::Result { + let mut result = None; + self.fetch_range_batch(std::slice::from_ref(&range), |_, response| { + result = Some(response); + Ok(()) + })?; + Ok(result.expect("one requested range")) } - fn fetch_range_batch(&self, ranges: &[Range]) -> io::Result> { - debug_assert!(ranges.len() <= RANGE_READ_CONCURRENCY); + fn fetch_range_batch( + &self, + ranges: &[Range], + mut consume: impl FnMut(usize, RangeResponse) -> io::Result<()>, + ) -> io::Result<()> { let reader = Arc::clone(&self.reader); - let permits = Arc::clone(&self.permits); + let io_permits = Arc::clone(&self.limiter.io_permits); + let response_permits = Arc::clone(&self.limiter.response_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()); + let io_limit = self.limiter.io_limit; + let response_limit = self.limiter.response_limit; + let (sender, mut receiver) = tokio::sync::mpsc::channel(io_limit); self.runtime.spawn(async move { - let fetched = try_join_all(requested.iter().cloned().map(|range| { + let fetched = stream::iter(requested.into_iter().enumerate().map(|(index, range)| { let reader = Arc::clone(&reader); - let permits = Arc::clone(&permits); + let io_permits = Arc::clone(&io_permits); + let response_permits = Arc::clone(&response_permits); + let sender = sender.clone(); let path = path.clone(); let stats = stats.clone(); async move { + let response_permit = response_permits + .acquire_owned() + .await + .map_err(|_| { + io::Error::other("vindex range response limiter closed") + })?; let permit_wait_start = stats.as_ref().map(|_| Instant::now()); - let permit = permits.acquire_owned().await; + let permit = io_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(|_| { + let io_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 { + let in_flight_read = stats.as_ref().map(|stats| { stats.file_read_calls.fetch_add(1, Ordering::Relaxed); - } - let data = reader.read(range.clone()).await.map_err(|error| { + let active = stats.in_flight_reads.fetch_add(1, Ordering::Relaxed) + 1; + stats + .peak_in_flight_reads + .fetch_max(active, Ordering::Relaxed); + InFlightRead(&stats.in_flight_reads) + }); + let read_result = reader.read(range.clone()).await; + drop(in_flight_read); + drop(io_permit); + let data = read_result.map_err(|error| { io::Error::other(format!( "failed to read vindex file '{path}' range {}..{}: {error}", range.start, range.end @@ -257,24 +337,49 @@ impl VindexFileReader { .returned_bytes .fetch_add(data.len() as u64, Ordering::Relaxed); } - Ok(data) + sender + .send(Ok((index, data, response_permit))) + .await + .map_err(|_| io::Error::other("vindex range read receiver closed")) } })) - .await; - let _ = sender.send(fetched); + .buffer_unordered(response_limit); + futures::pin_mut!(fetched); + while let Some(result) = fetched.next().await { + if let Err(error) = result { + let _ = sender.send(Err(error)).await; + return; + } + } + }); + + let mut io_wait_nanos = 0u64; + let result = (0..ranges.len()).try_for_each(|_| { + let wait_start = self.stats.as_ref().map(|_| Instant::now()); + let fetched = receiver.blocking_recv(); + if let Some(start) = wait_start { + io_wait_nanos = io_wait_nanos.saturating_add(start.elapsed().as_nanos() as u64); + } + let (index, data, response_permit) = fetched.ok_or_else(|| { + io::Error::other(format!( + "vindex range read task for '{}' was cancelled", + self.path + )) + })??; + consume( + index, + RangeResponse { + data, + _permit: response_permit, + }, + ) }); - let result = receiver.recv(); - if let (Some(stats), Some(start)) = (&self.stats, wait_start) { + if let Some(stats) = &self.stats { stats .io_wait_nanos - .fetch_add(start.elapsed().as_nanos() as u64, Ordering::Relaxed); + .fetch_add(io_wait_nanos, Ordering::Relaxed); } - result.map_err(|_| { - io::Error::other(format!( - "vindex range read task for '{}' was cancelled", - self.path - )) - })? + result } fn read_many(&self, requests: &mut [ReadRequest<'_>]) -> io::Result<()> { @@ -319,20 +424,42 @@ impl VindexFileReader { }); } - for batch in merged.chunks(RANGE_READ_CONCURRENCY) { - let ranges: Vec<_> = batch.iter().map(|merged| merged.range.clone()).collect(); - let fetched = self.fetch_range_batch(&ranges)?; - for (merged_range, data) in batch.iter().zip(fetched) { - for &request_index in &merged_range.request_indices { - let request = &mut requests[request_index]; - let start = (request.pos - merged_range.range.start) as usize; - request - .buf - .copy_from_slice(&data[start..start + request.buf.len()]); - } - } + if let Some(stats) = &self.stats { + let chunk_size = merged.len() as u64; + stats + .read_many_merged_ranges + .fetch_add(chunk_size, Ordering::Relaxed); + stats.read_many_chunks.fetch_add(1, Ordering::Relaxed); + stats + .read_many_chunk_size_sum + .fetch_add(chunk_size, Ordering::Relaxed); + let _ = stats.read_many_chunk_size_min.fetch_update( + Ordering::Relaxed, + Ordering::Relaxed, + |current| { + Some(if current == 0 { + chunk_size + } else { + current.min(chunk_size) + }) + }, + ); + stats + .read_many_chunk_size_max + .fetch_max(chunk_size, Ordering::Relaxed); } - Ok(()) + let ranges: Vec<_> = merged.iter().map(|merged| merged.range.clone()).collect(); + self.fetch_range_batch(&ranges, |merged_index, response| { + let merged_range = &merged[merged_index]; + for &request_index in &merged_range.request_indices { + let request = &mut requests[request_index]; + let start = (request.pos - merged_range.range.start) as usize; + request + .buf + .copy_from_slice(&response.data[start..start + request.buf.len()]); + } + Ok(()) + }) } } @@ -371,7 +498,7 @@ impl SeekRead for VindexFileReader { Ok(Some(Self { reader: Arc::clone(&self.reader), runtime: self.runtime.clone(), - permits: Arc::clone(&self.permits), + limiter: self.limiter.clone(), file_size: self.file_size, path: self.path.clone(), scalar_cache: None, @@ -380,9 +507,8 @@ impl SeekRead for VindexFileReader { } fn read_capabilities(&self) -> SeekReadCapabilities { - // This adapter accepts any number of ranges and splits them internally. - // The efficient window size depends on the underlying FileRead backend, - // so leave both storage-specific hints unspecified. + // `max_ranges_per_pread` is a planning hint, not an I/O concurrency limit. + // This adapter accepts any number of ranges and limits concurrent I/O with permits. SeekReadCapabilities::default() } } @@ -396,6 +522,14 @@ mod tests { use std::sync::Mutex; use std::time::Duration; + async fn acquire_test_permits(semaphore: &tokio::sync::Semaphore, permits: u32, message: &str) { + tokio::time::timeout(Duration::from_secs(5), semaphore.acquire_many(permits)) + .await + .expect(message) + .unwrap() + .forget(); + } + struct TrackingRead { data: Bytes, ranges: Mutex>>, @@ -433,6 +567,47 @@ mod tests { data: Bytes, active: AtomicUsize, max_active: AtomicUsize, + started: tokio::sync::Semaphore, + release: tokio::sync::Semaphore, + } + + struct FailOnceRead { + data: Bytes, + calls: AtomicUsize, + } + + struct DropTrackedPayload { + data: Vec, + dropped: Arc, + } + + impl AsRef<[u8]> for DropTrackedPayload { + fn as_ref(&self) -> &[u8] { + &self.data + } + } + + impl Drop for DropTrackedPayload { + fn drop(&mut self) { + self.dropped.add_permits(1); + } + } + + struct StreamingTrackingRead { + stride: u64, + first_started: tokio::sync::Semaphore, + release_first: tokio::sync::Semaphore, + dropped: Arc, + } + + struct BenchmarkRead { + data: Bytes, + calls: AtomicUsize, + active: AtomicUsize, + max_active: AtomicUsize, + fast_delay: Duration, + slow_every: usize, + slow_delay: Duration, } #[async_trait] @@ -448,12 +623,46 @@ mod tests { async fn read(&self, range: Range) -> crate::Result { let active = self.active.fetch_add(1, Ordering::SeqCst) + 1; self.max_active.fetch_max(active, Ordering::SeqCst); - tokio::time::sleep(Duration::from_millis(25)).await; + self.started.add_permits(1); + self.release.acquire().await.unwrap().forget(); self.active.fetch_sub(1, Ordering::SeqCst); Ok(self.data.slice(range.start as usize..range.end as usize)) } } + #[async_trait] + impl FileRead for FailOnceRead { + async fn read(&self, range: Range) -> crate::Result { + if self.calls.fetch_add(1, Ordering::SeqCst) == 0 { + return Err(crate::Error::UnexpectedError { + message: "injected range read failure".to_string(), + source: None, + }); + } + Ok(self.data.slice(range.start as usize..range.end as usize)) + } + } + + #[async_trait] + impl FileRead for StreamingTrackingRead { + async fn read(&self, range: Range) -> crate::Result { + if range.start == 0 { + self.first_started.add_permits(1); + acquire_test_permits( + &self.release_first, + 1, + "test did not release the first range read", + ) + .await; + } + let value = (range.start / self.stride + 1) as u8; + Ok(Bytes::from_owner(DropTrackedPayload { + data: vec![value; (range.end - range.start) as usize], + dropped: Arc::clone(&self.dropped), + })) + } + } + #[async_trait] impl FileRead for TrackingRead { async fn read(&self, range: Range) -> crate::Result { @@ -466,6 +675,113 @@ mod tests { } } + #[async_trait] + impl FileRead for BenchmarkRead { + async fn read(&self, range: Range) -> crate::Result { + let call = self.calls.fetch_add(1, Ordering::Relaxed) + 1; + let active = self.active.fetch_add(1, Ordering::Relaxed) + 1; + self.max_active.fetch_max(active, Ordering::Relaxed); + let delay = if self.slow_every != 0 && call.is_multiple_of(self.slow_every) { + self.slow_delay + } else { + self.fast_delay + }; + if !delay.is_zero() { + tokio::time::sleep(delay).await; + } + self.active.fetch_sub(1, Ordering::Relaxed); + Ok(self.data.slice(range.start as usize..range.end as usize)) + } + } + + async fn run_range_read_benchmark( + case: &str, + range_count: usize, + iterations: usize, + fast_delay: Duration, + slow_every: usize, + slow_delay: Duration, + ) { + const RANGE_SIZE: usize = 4 * 1024; + let concurrency = 64; + let stride = RANGE_COALESCE_GAP as usize + RANGE_SIZE + 1; + let source = Arc::new(BenchmarkRead { + data: Bytes::from(vec![7u8; range_count * stride]), + calls: AtomicUsize::new(0), + active: AtomicUsize::new(0), + max_active: AtomicUsize::new(0), + fast_delay, + slow_every, + slow_delay, + }); + let reader_source: Arc = source.clone(); + let mut reader = VindexFileReader::new_with_limiter( + reader_source, + tokio::runtime::Handle::current(), + RangeReadLimiter::new(concurrency), + source.data.len() as u64, + "benchmark-index".to_string(), + ); + let task_source = source.clone(); + let elapsed = tokio::task::spawn_blocking(move || { + let mut buffers = vec![[0u8; RANGE_SIZE]; range_count]; + let mut run_iteration = || { + let mut requests = buffers + .iter_mut() + .enumerate() + .map(|(index, buffer)| { + ReadRequest::new((index * stride) as u64, buffer.as_mut_slice()) + }) + .collect::>(); + reader.pread(&mut requests).unwrap(); + std::hint::black_box(&buffers); + }; + + run_iteration(); + task_source.calls.store(0, Ordering::Relaxed); + task_source.max_active.store(0, Ordering::Relaxed); + let start = Instant::now(); + for _ in 0..iterations { + run_iteration(); + } + start.elapsed() + }) + .await + .unwrap(); + + let total_ranges = range_count * iterations; + let total_bytes = total_ranges * RANGE_SIZE; + let ranges_per_second = total_ranges as f64 / elapsed.as_secs_f64(); + let mib_per_second = total_bytes as f64 / (1024.0 * 1024.0) / elapsed.as_secs_f64(); + let peak_in_flight = source.max_active.load(Ordering::Relaxed); + assert_eq!(source.calls.load(Ordering::Relaxed), total_ranges); + assert!(peak_in_flight <= concurrency); + eprintln!( + "vindex_range_read_benchmark case={case} profile={} concurrency={concurrency} range_bytes={RANGE_SIZE} ranges_per_iteration={range_count} iterations={iterations} elapsed_ms={:.3} ranges_per_second={ranges_per_second:.0} mib_per_second={mib_per_second:.2} peak_in_flight={peak_in_flight}", + if cfg!(debug_assertions) { + "debug" + } else { + "release" + }, + elapsed.as_secs_f64() * 1000.0, + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[ignore = "manual performance comparison; run with --release --ignored --nocapture"] + async fn vindex_range_read_benchmark() { + run_range_read_benchmark("hot_cache", 1024, 50, Duration::ZERO, 0, Duration::ZERO).await; + run_range_read_benchmark( + "oss_straggler", + 256, + 10, + Duration::from_millis(1), + 64, + Duration::from_millis(10), + ) + .await; + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn scalar_reads_reuse_bounded_read_ahead() { let data = Bytes::from((0..200_000).map(|value| value as u8).collect::>()); @@ -714,8 +1030,14 @@ mod tests { stats.file_read_calls, stats.returned_bytes, stats.read_ahead_hits, + stats.peak_in_flight_reads, + stats.read_many_merged_ranges, + stats.read_many_chunks, + stats.read_many_chunk_size_sum, + stats.read_many_chunk_size_min, + stats.read_many_chunk_size_max, ), - (2, 256, 2, 256, 0) + (2, 256, 2, 256, 0, 1, 0, 0, 0, 0, 0) ); assert!(stats.io_wait_nanos > 0); } @@ -744,43 +1066,68 @@ mod tests { ReadRequest::new(20_000, &mut third), ]) .unwrap(); + let mut fourth = [0u8; 4]; + let mut fifth = [0u8; 4]; + reader + .pread(&mut [ + ReadRequest::new(40, &mut fourth), + ReadRequest::new(48, &mut fifth), + ]) + .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); + assert_eq!( + ( + stats.logical_ranges, + stats.requested_bytes, + stats.file_read_calls, + stats.returned_bytes, + stats.read_ahead_hits, + stats.peak_in_flight_reads, + stats.read_many_merged_ranges, + stats.read_many_chunks, + stats.read_many_chunk_size_sum, + stats.read_many_chunk_size_min, + stats.read_many_chunk_size_max, + ), + (5, 20, 3, 28, 0, 1, 3, 2, 3, 1, 2) + ); + assert!(stats.io_wait_nanos > 0); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn shared_permits_bound_reads_across_independent_readers() { + async fn clones_share_range_read_permits() { let data = Bytes::from(vec![8u8; 1024]); let tracking = Arc::new(ConcurrencyTrackingRead { data: data.clone(), active: AtomicUsize::new(0), max_active: AtomicUsize::new(0), + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), }); - let permits = Arc::new(tokio::sync::Semaphore::new(0)); - let make_reader = |path: &str| { - let source: Arc = tracking.clone(); - VindexFileReader::new_with_permits( - source, - tokio::runtime::Handle::current(), - Arc::clone(&permits), - data.len() as u64, - path.to_string(), - ) - }; - let mut first_reader = make_reader("first.index"); - let mut second_reader = make_reader("second.index"); + let limiter = RangeReadLimiter::new(1); + let source: Arc = tracking.clone(); + let mut first_reader = VindexFileReader::new_with_limiter( + source, + tokio::runtime::Handle::current(), + limiter, + data.len() as u64, + "index".to_string(), + ); 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 mut second_reader = first_reader.try_clone_reader().unwrap().unwrap(); + assert!(Arc::ptr_eq( + &first_reader.limiter.io_permits, + &second_reader.limiter.io_permits + )); + assert!(Arc::ptr_eq( + &first_reader.limiter.response_permits, + &second_reader.limiter.response_permits + )); let first = tokio::task::spawn_blocking(move || { let mut output = [0u8; 128]; @@ -794,8 +1141,12 @@ mod tests { .pread(&mut [ReadRequest::new(128, &mut output)]) .unwrap(); }); - tokio::time::sleep(Duration::from_millis(10)).await; - permits.add_permits(1); + tracking.started.acquire().await.unwrap().forget(); + assert_eq!(tracking.max_active.load(Ordering::SeqCst), 1); + tracking.release.add_permits(1); + tracking.started.acquire().await.unwrap().forget(); + assert_eq!(tracking.max_active.load(Ordering::SeqCst), 1); + tracking.release.add_permits(1); first.await.unwrap(); second.await.unwrap(); @@ -805,25 +1156,45 @@ mod tests { assert!(stats.io_wait_nanos >= stats.range_permit_wait_nanos); } + #[test] + fn response_limit_saturates_at_semaphore_max_permits() { + let io_limit = tokio::sync::Semaphore::MAX_PERMITS / 2 + 1; + let limiter = RangeReadLimiter::new(io_limit); + + assert_eq!(limiter.io_limit, io_limit); + assert_eq!(limiter.response_limit, tokio::sync::Semaphore::MAX_PERMITS); + assert_eq!(limiter.io_permits.available_permits(), io_limit); + assert_eq!( + limiter.response_permits.available_permits(), + tokio::sync::Semaphore::MAX_PERMITS + ); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn read_many_caps_each_batch_at_32_ranges() { - let range_count = RANGE_READ_CONCURRENCY + 1; + async fn configured_range_read_concurrency_can_exceed_32() { + let configured_concurrency = 64; + let range_count = configured_concurrency + 1; let stride = RANGE_COALESCE_GAP + 2; let data = Bytes::from(vec![8u8; range_count * stride as usize]); let tracking = Arc::new(ConcurrencyTrackingRead { data: data.clone(), active: AtomicUsize::new(0), max_active: AtomicUsize::new(0), + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), }); let source: Arc = tracking.clone(); - let mut reader = VindexFileReader::new( + let mut reader = VindexFileReader::new_with_limiter( source, tokio::runtime::Handle::current(), + RangeReadLimiter::new(configured_concurrency), 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 read = tokio::task::spawn_blocking(move || { let mut buffers = vec![[0u8; 1]; range_count]; let mut requests = buffers .iter_mut() @@ -831,14 +1202,292 @@ mod tests { .map(|(index, buffer)| ReadRequest::new(index as u64 * stride, buffer)) .collect::>(); reader.pread(&mut requests).unwrap(); - }) - .await - .unwrap(); + }); + tokio::time::timeout( + Duration::from_secs(5), + tracking.started.acquire_many(configured_concurrency as u32), + ) + .await + .expect("configured range reads did not start") + .unwrap() + .forget(); + assert_eq!(tracking.max_active.load(Ordering::SeqCst), 64); + tracking.release.add_permits(configured_concurrency); + tracking.started.acquire().await.unwrap().forget(); + tracking.release.add_permits(1); + read.await.unwrap(); + + assert_eq!(tracking.max_active.load(Ordering::SeqCst), 64); + assert!(tracking.max_active.load(Ordering::SeqCst) > 32); + let snapshot = stats.snapshot(); assert_eq!( - tracking.max_active.load(Ordering::SeqCst), - RANGE_READ_CONCURRENCY + ( + snapshot.logical_ranges, + snapshot.requested_bytes, + snapshot.file_read_calls, + snapshot.returned_bytes, + snapshot.read_ahead_hits, + snapshot.peak_in_flight_reads, + snapshot.read_many_merged_ranges, + snapshot.read_many_chunks, + snapshot.read_many_chunk_size_sum, + snapshot.read_many_chunk_size_min, + snapshot.read_many_chunk_size_max, + ), + (65, 65, 65, 65, 0, 64, 65, 1, 65, 65, 65) + ); + assert!(snapshot.io_wait_nanos > 0); + assert!(snapshot.range_permit_wait_nanos > 0); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn range_reads_refill_before_slowest_batch_member_finishes() { + let concurrency = 2; + let range_count = 3; + let stride = RANGE_COALESCE_GAP + 2; + let data = Bytes::from(vec![8u8; range_count * stride as usize]); + let tracking = Arc::new(ConcurrencyTrackingRead { + data: data.clone(), + active: AtomicUsize::new(0), + max_active: AtomicUsize::new(0), + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), + }); + let source: Arc = tracking.clone(); + let mut reader = VindexFileReader::new_with_limiter( + source, + tokio::runtime::Handle::current(), + RangeReadLimiter::new(concurrency), + data.len() as u64, + "index".to_string(), + ); + + let read = tokio::task::spawn_blocking(move || { + let mut buffers = vec![[0u8; 1]; range_count]; + let mut requests = buffers + .iter_mut() + .enumerate() + .map(|(index, buffer)| ReadRequest::new(index as u64 * stride, buffer)) + .collect::>(); + reader.pread(&mut requests).unwrap(); + }); + + tracking + .started + .acquire_many(concurrency as u32) + .await + .unwrap() + .forget(); + tracking.release.add_permits(1); + tokio::time::timeout(Duration::from_secs(1), tracking.started.acquire()) + .await + .expect("next range did not refill while another range was still running") + .unwrap() + .forget(); + tracking.release.add_permits(concurrency); + read.await.unwrap(); + + assert_eq!(tracking.max_active.load(Ordering::SeqCst), concurrency); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn response_copy_and_buffers_are_bounded_across_clones() { + let stride = RANGE_COALESCE_GAP + 2; + let data = Bytes::from(vec![8u8; 3 * stride as usize]); + let tracking = Arc::new(ConcurrencyTrackingRead { + data: data.clone(), + active: AtomicUsize::new(0), + max_active: AtomicUsize::new(0), + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), + }); + let source: Arc = tracking.clone(); + let first_reader = VindexFileReader::new_with_limiter( + source, + tokio::runtime::Handle::current(), + RangeReadLimiter::new(1), + data.len() as u64, + "index".to_string(), + ); + let second_reader = first_reader.try_clone_reader().unwrap().unwrap(); + let (copy_started_tx, copy_started_rx) = std::sync::mpsc::channel(); + let (release_copy_tx, release_copy_rx) = std::sync::mpsc::channel(); + + let first = tokio::task::spawn_blocking(move || { + first_reader + .fetch_range_batch(&[0..1, stride..stride + 1], |index, _| { + if index == 0 { + copy_started_tx.send(()).unwrap(); + release_copy_rx.recv().unwrap(); + } + Ok(()) + }) + .unwrap(); + }); + + acquire_test_permits(&tracking.started, 1, "first range read did not start").await; + tracking.release.add_permits(1); + copy_started_rx + .recv_timeout(Duration::from_secs(5)) + .expect("first response did not reach the copy stage"); + acquire_test_permits(&tracking.started, 1, "second range read did not start").await; + let second = tokio::task::spawn_blocking(move || { + let range = 2 * stride..2 * stride + 1; + second_reader + .fetch_range_batch(std::slice::from_ref(&range), |_, _| Ok(())) + .unwrap(); + }); + tracking.release.add_permits(1); + let third_started_early = + tokio::time::timeout(Duration::from_secs(1), tracking.started.acquire()) + .await + .map(|permit| permit.unwrap().forget()) + .is_ok(); + release_copy_tx.send(()).unwrap(); + if !third_started_early { + acquire_test_permits(&tracking.started, 1, "third range read did not start").await; + } + tracking.release.add_permits(1); + first.await.unwrap(); + second.await.unwrap(); + + assert!( + !third_started_early, + "more than 2x the I/O concurrency was retained as responses" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn single_range_response_holds_permit_until_consumed() { + let data = Bytes::from(vec![8u8; 384]); + let limiter = RangeReadLimiter::new(1); + let response_permits = Arc::clone(&limiter.response_permits); + let source: Arc = TrackingRead::new(data.clone()); + let first_reader = VindexFileReader::new_with_limiter( + source, + tokio::runtime::Handle::current(), + limiter, + data.len() as u64, + "index".to_string(), + ); + let second_reader = first_reader.try_clone_reader().unwrap().unwrap(); + let third_reader = first_reader.try_clone_reader().unwrap().unwrap(); + + let first = tokio::task::spawn_blocking(move || first_reader.fetch_exact(0..128).unwrap()) + .await + .unwrap(); + let second = + tokio::task::spawn_blocking(move || second_reader.fetch_exact(128..256).unwrap()) + .await + .unwrap(); + let mut third = + tokio::task::spawn_blocking(move || third_reader.fetch_exact(256..384).unwrap()); + + assert!(tokio::time::timeout(Duration::from_secs(1), &mut third) + .await + .is_err()); + drop(first); + let third = tokio::time::timeout(Duration::from_secs(5), third) + .await + .expect("third single-range response did not start after a permit was released") + .unwrap(); + drop(second); + drop(third); + assert_eq!(response_permits.available_permits(), 2); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn completed_range_buffers_are_released_before_slowest_read() { + let concurrency = 2; + let range_count = 3; + let stride = RANGE_COALESCE_GAP + 2; + let dropped = Arc::new(tokio::sync::Semaphore::new(0)); + let tracking = Arc::new(StreamingTrackingRead { + stride, + first_started: tokio::sync::Semaphore::new(0), + release_first: tokio::sync::Semaphore::new(0), + dropped: Arc::clone(&dropped), + }); + let source: Arc = tracking.clone(); + let mut reader = VindexFileReader::new_with_limiter( + source, + tokio::runtime::Handle::current(), + RangeReadLimiter::new(concurrency), + range_count as u64 * stride, + "index".to_string(), ); + + let read = tokio::task::spawn_blocking(move || { + let mut buffers = vec![[0u8; 1]; range_count]; + let mut requests = buffers + .iter_mut() + .enumerate() + .map(|(index, buffer)| ReadRequest::new(index as u64 * stride, buffer)) + .collect::>(); + reader.pread(&mut requests).unwrap(); + buffers + }); + + acquire_test_permits(&tracking.first_started, 1, "first range read did not start").await; + let released = tokio::time::timeout(Duration::from_secs(5), dropped.acquire_many(2)).await; + tracking.release_first.add_permits(1); + released + .expect("completed range buffers were retained by the slowest read") + .unwrap() + .forget(); + + let output = tokio::time::timeout(Duration::from_secs(5), read) + .await + .expect("range read task did not finish") + .unwrap(); + assert_eq!(output, vec![[1], [2], [3]]); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn failed_range_read_releases_permit() { + let data = Bytes::from(vec![9u8; 1024]); + let source: Arc = Arc::new(FailOnceRead { + data: data.clone(), + calls: AtomicUsize::new(0), + }); + let limiter = RangeReadLimiter::new(1); + let io_permits = Arc::clone(&limiter.io_permits); + let response_permits = Arc::clone(&limiter.response_permits); + let mut reader = VindexFileReader::new_with_limiter( + source, + tokio::runtime::Handle::current(), + limiter, + data.len() as u64, + "index".to_string(), + ); + + let (mut reader, error) = tokio::task::spawn_blocking(move || { + let mut output = [0u8; 128]; + let error = reader + .pread(&mut [ReadRequest::new(0, &mut output)]) + .unwrap_err(); + (reader, error) + }) + .await + .unwrap(); + assert_eq!(error.kind(), io::ErrorKind::Other); + assert_eq!(io_permits.available_permits(), 1); + assert_eq!(response_permits.available_permits(), 2); + + tokio::time::timeout( + Duration::from_secs(5), + tokio::task::spawn_blocking(move || { + let mut output = [0u8; 128]; + reader + .pread(&mut [ReadRequest::new(0, &mut output)]) + .unwrap(); + assert_eq!(output, [9u8; 128]); + }), + ) + .await + .expect("range read blocked after an error") + .unwrap(); } #[test] diff --git a/docs/src/sql.md b/docs/src/sql.md index eaf9758a..1440cb67 100644 --- a/docs/src/sql.md +++ b/docs/src/sql.md @@ -2064,9 +2064,15 @@ deletion vectors enabled. | `btree-index.fallback-scan-max-size` | `256mb` | Maximum total size of selected BTree global-index files for fallback scans used by range/between and suffix/contains/complex LIKE predicates; `0` disables BTree fallback index scans. | | `bitmap-index.fallback-scan-max-size` | `256mb` | Maximum total size of selected bitmap global-index files for fallback scans used by range/between and suffix/contains/complex LIKE predicates; `0` disables bitmap fallback index scans. | | `global-index.search-mode` | `fast` | Global index coverage mode for reads: `fast`, `full`, or `detail`. | -| `global-index.thread-num` | `32` | Number of threads used to search global index fields concurrently; must be greater than 0 and must not exceed the runtime's task limit. | +| `global-index.thread-num` | `32` | Number of concurrent global-index search tasks; must be greater than 0 and must not exceed the runtime's task limit. This does not limit Vindex file range reads. | +| `global-index.vindex.read-thread-num` | `64` | Maximum number of concurrent Vindex file range reads shared by one search operation; must be greater than 0 and must not exceed the runtime semaphore limit. | | `global-index.column-update-action` | `THROW_ERROR` | What a commit does when it updates an indexed column: `THROW_ERROR` rejects the commit, `DROP_PARTITION_INDEX` drops the affected partition index instead. | +`global-index.vindex.read-thread-num` is independent of `global-index.thread-num`. +When upgrading a table that sets `global-index.thread-num`, set the new option +explicitly to the same value if Vindex range reads should keep the previous limit; +otherwise they use the new default of `64`. + ### Variant Shredding Options Set these as table options when writing `VARIANT` columns to Parquet. The