From 6f86d0c7342b2d42e34903cd5132846799e8a032 Mon Sep 17 00:00:00 2001 From: yantian Date: Mon, 17 Aug 2026 10:04:06 +0800 Subject: [PATCH 01/14] perf(vindex): configure range read concurrency --- crates/paimon/src/spec/core_options.rs | 76 ++++++ .../paimon/src/table/vector_search_builder.rs | 93 +++++-- crates/paimon/src/vindex/range_reader.rs | 240 +++++++++++++++--- 3 files changed, 344 insertions(+), 65 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 1d45d49a..fd92b75c 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_RANGE_READ_THREAD_NUM_OPTION: &str = "global-index.range-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_RANGE_READ_THREAD_NUM: usize = 32; +const MAX_GLOBAL_INDEX_RANGE_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; @@ -728,6 +731,35 @@ impl<'a> CoreOptions<'a> { Ok(value as usize) } + /// Maximum number of concurrent range reads shared by Vindex readers in one + /// search operation. This is independent of [`Self::global_index_thread_num`]. + pub fn global_index_range_read_thread_num(&self) -> crate::Result { + let value = self + .parse_i64_option(GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION)? + .unwrap_or(DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM as i64); + if value <= 0 { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{}' must be greater than 0, got: {}", + GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION, value + ), + source: None, + }); + } + if value > MAX_GLOBAL_INDEX_RANGE_READ_THREAD_NUM { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{}' must not exceed {}, got: {}", + GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION, + MAX_GLOBAL_INDEX_RANGE_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 +1552,10 @@ mod tests { 100_000 ); assert_eq!(core_options.global_index_thread_num().unwrap(), 32); + assert_eq!( + core_options.global_index_range_read_thread_num().unwrap(), + 32 + ); assert_eq!( core_options.sorted_index_records_per_range().unwrap(), 100_000 @@ -1764,6 +1800,46 @@ mod tests { ); } + #[test] + fn test_global_index_range_read_thread_num_default_and_custom() { + assert_eq!( + CoreOptions::new(&HashMap::new()) + .global_index_range_read_thread_num() + .unwrap(), + 32 + ); + + for value in [32, 64] { + let options = HashMap::from([( + GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION.to_string(), + value.to_string(), + )]); + assert_eq!( + CoreOptions::new(&options) + .global_index_range_read_thread_num() + .unwrap(), + value + ); + } + } + + #[test] + fn test_global_index_range_read_thread_num_rejects_invalid_values() { + for value in [ + "0".to_string(), + "abc".to_string(), + (MAX_GLOBAL_INDEX_RANGE_READ_THREAD_NUM + 1).to_string(), + ] { + let options = + HashMap::from([(GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION.to_string(), value)]); + let err = CoreOptions::new(&options) + .global_index_range_read_thread_num() + .expect_err("invalid range-read-thread-num should fail"); + assert!(matches!(err, crate::Error::DataInvalid { message, .. } + if message.contains(GLOBAL_INDEX_RANGE_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..ac3a813a 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -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_sizes={:?}", file, query_count, stats.logical_ranges, @@ -154,9 +154,24 @@ 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_sizes, ); } +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_range_read_thread_num()?, + )) +} + pub struct VectorSearchBuilder<'a> { table: &'a Table, vector_column: Option, @@ -844,15 +859,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 +894,16 @@ 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_permits = match backend { + VectorIndexBackend::Vindex => Some(Arc::new(tokio::sync::Semaphore::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_permits = loader_range_read_permits.clone(); let path = segment.path.clone(); let file_size = segment.file_size; Box::pin(async move { @@ -911,7 +932,8 @@ async fn plan_and_search_pk_candidates_batch( VindexFileReader::new_with_permits( Arc::new(file_reader), current_tokio_runtime_handle()?, - range_read_permits, + range_read_permits.expect("Vindex range-read permits"), + range_read_concurrency, file_size, path, ), @@ -1629,18 +1651,28 @@ 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_concurrency, range_read_permits) = + if vindex_entry_count == 0 { + (1, 0, None) + } else { + let (index_parallelism, range_read_concurrency) = + vindex_concurrency_limits(&core_options, vindex_entry_count, concurrency)?; + ( + index_parallelism, + range_read_concurrency, + Some(Arc::new(tokio::sync::Semaphore::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_permits = range_read_permits.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"); @@ -1724,7 +1756,8 @@ async fn evaluate_batch_vector_search( let source = VindexFileReader::new_with_permits( Arc::new(file_reader), runtime, - range_read_permits, + range_read_permits.expect("Vindex range-read permits"), + range_read_concurrency, file_size, file_name.clone(), ); @@ -3446,11 +3479,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, 32) + ); + assert_eq!( + vindex_concurrency_limits(&default_core, 8, 4).unwrap(), + (4, 32) + ); + + let options = HashMap::from([( + "global-index.range-read-thread-num".to_string(), + "64".to_string(), + )]); + let core = CoreOptions::new(&options); + assert_eq!(vindex_concurrency_limits(&core, 1, 32).unwrap(), (1, 64)); + assert_eq!(vindex_concurrency_limits(&core, 8, 4).unwrap(), (4, 64)); } fn make_field(id: i32, name: &str) -> DataField { diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 0f5d965c..517ad38f 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -16,6 +16,8 @@ // under the License. use crate::io::FileRead; +#[cfg(test)] +use crate::spec::DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM; use crate::vindex::vector_search_timing_enabled; use bytes::Bytes; use futures::future::try_join_all; @@ -29,7 +31,6 @@ 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, @@ -66,9 +67,13 @@ pub(crate) struct RangeIoStats { read_ahead_hits: AtomicU64, io_wait_nanos: AtomicU64, range_permit_wait_nanos: AtomicU64, + peak_in_flight_reads: AtomicU64, + read_many_merged_ranges: AtomicU64, + read_many_chunks: AtomicU64, + read_many_chunk_sizes: std::sync::Mutex>, } -#[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 +82,10 @@ 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_sizes: Vec, } impl RangeIoStats { @@ -89,6 +98,10 @@ 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_sizes: self.read_many_chunk_sizes.lock().unwrap().clone(), } } } @@ -100,6 +113,7 @@ pub(crate) struct VindexFileReader { reader: Arc, runtime: tokio::runtime::Handle, permits: Arc, + max_range_read_concurrency: usize, file_size: u64, path: String, scalar_cache: Option, @@ -117,7 +131,10 @@ impl VindexFileReader { Self::new_with_permits( reader, runtime, - Arc::new(tokio::sync::Semaphore::new(RANGE_READ_CONCURRENCY)), + Arc::new(tokio::sync::Semaphore::new( + DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM, + )), + DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM, file_size, path, ) @@ -127,6 +144,7 @@ impl VindexFileReader { reader: Arc, runtime: tokio::runtime::Handle, permits: Arc, + max_range_read_concurrency: usize, file_size: u64, path: String, ) -> Self { @@ -134,6 +152,7 @@ impl VindexFileReader { reader, runtime, permits, + max_range_read_concurrency, file_size, path, scalar_cache: None, @@ -206,12 +225,13 @@ impl VindexFileReader { } fn fetch_range_batch(&self, ranges: &[Range]) -> io::Result> { - debug_assert!(ranges.len() <= RANGE_READ_CONCURRENCY); + debug_assert!(ranges.len() <= self.max_range_read_concurrency); let reader = Arc::clone(&self.reader); let permits = Arc::clone(&self.permits); let path = self.path.clone(); let requested = ranges.to_vec(); let stats = self.stats.clone(); + let max_range_read_concurrency = self.max_range_read_concurrency; let (sender, receiver) = mpsc::sync_channel(1); let wait_start = self.stats.as_ref().map(|_| Instant::now()); self.runtime.spawn(async move { @@ -222,7 +242,7 @@ impl VindexFileReader { let stats = stats.clone(); async move { let permit_wait_start = stats.as_ref().map(|_| Instant::now()); - let permit = permits.acquire_owned().await; + let permit = Arc::clone(&permits).acquire_owned().await; if let (Some(stats), Some(start)) = (&stats, permit_wait_start) { stats .range_permit_wait_nanos @@ -234,6 +254,12 @@ impl VindexFileReader { let expected = (range.end - range.start) as usize; if let Some(stats) = &stats { stats.file_read_calls.fetch_add(1, Ordering::Relaxed); + stats.peak_in_flight_reads.fetch_max( + max_range_read_concurrency + .saturating_sub(permits.available_permits()) + as u64, + Ordering::Relaxed, + ); } let data = reader.read(range.clone()).await.map_err(|error| { io::Error::other(format!( @@ -319,7 +345,20 @@ impl VindexFileReader { }); } - for batch in merged.chunks(RANGE_READ_CONCURRENCY) { + if let Some(stats) = &self.stats { + stats + .read_many_merged_ranges + .fetch_add(merged.len() as u64, Ordering::Relaxed); + } + for batch in merged.chunks(self.max_range_read_concurrency) { + if let Some(stats) = &self.stats { + stats.read_many_chunks.fetch_add(1, Ordering::Relaxed); + stats + .read_many_chunk_sizes + .lock() + .unwrap() + .push(batch.len()); + } 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) { @@ -372,6 +411,7 @@ impl SeekRead for VindexFileReader { reader: Arc::clone(&self.reader), runtime: self.runtime.clone(), permits: Arc::clone(&self.permits), + max_range_read_concurrency: self.max_range_read_concurrency, file_size: self.file_size, path: self.path.clone(), scalar_cache: None, @@ -380,9 +420,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 chunks them internally. SeekReadCapabilities::default() } } @@ -433,6 +472,13 @@ mod tests { data: Bytes, active: AtomicUsize, max_active: AtomicUsize, + started: tokio::sync::Semaphore, + release: tokio::sync::Semaphore, + } + + struct FailOnceRead { + data: Bytes, + calls: AtomicUsize, } #[async_trait] @@ -448,12 +494,26 @@ 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 TrackingRead { async fn read(&self, range: Range) -> crate::Result { @@ -714,8 +774,12 @@ 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_sizes, ), - (2, 256, 2, 256, 0) + (2, 256, 2, 256, 0, 1, 0, 0, Vec::new()) ); assert!(stats.io_wait_nanos > 0); } @@ -749,37 +813,46 @@ mod tests { .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_sizes, + ), + (3, 12, 2, 16, 0, 1, 2, 1, vec![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 permits = Arc::new(tokio::sync::Semaphore::new(1)); + let source: Arc = tracking.clone(); + let mut first_reader = VindexFileReader::new_with_permits( + source, + tokio::runtime::Handle::current(), + Arc::clone(&permits), + 1, + 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)); + let mut second_reader = first_reader.try_clone_reader().unwrap().unwrap(); assert!(Arc::ptr_eq(&first_reader.permits, &second_reader.permits)); let first = tokio::task::spawn_blocking(move || { @@ -794,8 +867,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(); @@ -806,24 +883,32 @@ mod tests { } #[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 permits = Arc::new(tokio::sync::Semaphore::new(configured_concurrency)); let source: Arc = tracking.clone(); - let mut reader = VindexFileReader::new( + let mut reader = VindexFileReader::new_with_permits( source, tokio::runtime::Handle::current(), + permits, + 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 +916,85 @@ 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_sizes, + ), + (65, 65, 65, 65, 0, 64, 65, 2, vec![64, 1]) + ); + assert!(snapshot.io_wait_nanos > 0); + assert!(snapshot.range_permit_wait_nanos > 0); + } + + #[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 permits = Arc::new(tokio::sync::Semaphore::new(1)); + let mut reader = VindexFileReader::new_with_permits( + source, + tokio::runtime::Handle::current(), + Arc::clone(&permits), + 1, + 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!(permits.available_permits(), 1); + + 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] From 3ec12dae675f1402b8b7795cd7f3a47b9c3361aa Mon Sep 17 00:00:00 2001 From: yantian Date: Mon, 17 Aug 2026 13:51:15 +0800 Subject: [PATCH 02/14] perf(vindex): remove range-read chunk barrier --- crates/paimon/src/vindex/range_reader.rs | 93 ++++++++++++++++++------ 1 file changed, 71 insertions(+), 22 deletions(-) diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 517ad38f..e7e16491 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -225,7 +225,6 @@ impl VindexFileReader { } fn fetch_range_batch(&self, ranges: &[Range]) -> io::Result> { - debug_assert!(ranges.len() <= self.max_range_read_concurrency); let reader = Arc::clone(&self.reader); let permits = Arc::clone(&self.permits); let path = self.path.clone(); @@ -350,25 +349,23 @@ impl VindexFileReader { .read_many_merged_ranges .fetch_add(merged.len() as u64, Ordering::Relaxed); } - for batch in merged.chunks(self.max_range_read_concurrency) { - if let Some(stats) = &self.stats { - stats.read_many_chunks.fetch_add(1, Ordering::Relaxed); - stats - .read_many_chunk_sizes - .lock() - .unwrap() - .push(batch.len()); - } - 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 { + stats.read_many_chunks.fetch_add(1, Ordering::Relaxed); + stats + .read_many_chunk_sizes + .lock() + .unwrap() + .push(merged.len()); + } + let ranges: Vec<_> = merged.iter().map(|merged| merged.range.clone()).collect(); + let fetched = self.fetch_range_batch(&ranges)?; + for (merged_range, data) in merged.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()]); } } Ok(()) @@ -421,7 +418,7 @@ impl SeekRead for VindexFileReader { fn read_capabilities(&self) -> SeekReadCapabilities { // `max_ranges_per_pread` is a planning hint, not an I/O concurrency limit. - // This adapter accepts any number of ranges and chunks them internally. + // This adapter accepts any number of ranges and limits concurrent I/O with permits. SeekReadCapabilities::default() } } @@ -947,12 +944,64 @@ mod tests { snapshot.read_many_chunks, snapshot.read_many_chunk_sizes, ), - (65, 65, 65, 65, 0, 64, 65, 2, vec![64, 1]) + (65, 65, 65, 65, 0, 64, 65, 1, vec![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 permits = Arc::new(tokio::sync::Semaphore::new(concurrency)); + let source: Arc = tracking.clone(); + let mut reader = VindexFileReader::new_with_permits( + source, + tokio::runtime::Handle::current(), + permits, + 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 failed_range_read_releases_permit() { let data = Bytes::from(vec![9u8; 1024]); From cbfdeea0b6fcc1a508617ad0fa3c81622ec581a7 Mon Sep 17 00:00:00 2001 From: yantian Date: Mon, 17 Aug 2026 17:33:31 +0800 Subject: [PATCH 03/14] fix(vindex): bound range read response memory --- crates/paimon/src/spec/core_options.rs | 13 +- crates/paimon/src/vindex/range_reader.rs | 290 +++++++++++++++++++++-- docs/src/sql.md | 8 +- 3 files changed, 279 insertions(+), 32 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index fd92b75c..05b94919 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -699,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_range_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)? @@ -732,7 +734,8 @@ impl<'a> CoreOptions<'a> { } /// Maximum number of concurrent range reads shared by Vindex readers in one - /// search operation. This is independent of [`Self::global_index_thread_num`]. + /// search operation (key `global-index.range-read-thread-num`, default 32). + /// This is independent of [`Self::global_index_thread_num`]. pub fn global_index_range_read_thread_num(&self) -> crate::Result { let value = self .parse_i64_option(GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION)? diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index e7e16491..95341526 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -20,12 +20,12 @@ use crate::io::FileRead; use crate::spec::DEFAULT_GLOBAL_INDEX_RANGE_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; @@ -118,6 +118,8 @@ pub(crate) struct VindexFileReader { path: String, scalar_cache: Option, stats: Option>, + #[cfg(test)] + permit_wait_pending: Option>, } impl VindexFileReader { @@ -157,6 +159,8 @@ impl VindexFileReader { path, scalar_cache: None, stats: vector_search_timing_enabled().then(|| Arc::new(RangeIoStats::default())), + #[cfg(test)] + permit_wait_pending: None, } } @@ -220,34 +224,63 @@ impl VindexFileReader { } 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")) + let mut result = None; + self.fetch_range_batch(std::slice::from_ref(&range), |_, data| { + result = Some(data); + Ok(()) + })?; + Ok(result.expect("one requested range")) } - fn fetch_range_batch(&self, ranges: &[Range]) -> io::Result> { + fn fetch_range_batch( + &self, + ranges: &[Range], + mut consume: impl FnMut(usize, Bytes) -> io::Result<()>, + ) -> io::Result<()> { let reader = Arc::clone(&self.reader); let permits = Arc::clone(&self.permits); let path = self.path.clone(); let requested = ranges.to_vec(); let stats = self.stats.clone(); let max_range_read_concurrency = self.max_range_read_concurrency; - let (sender, receiver) = mpsc::sync_channel(1); - let wait_start = self.stats.as_ref().map(|_| Instant::now()); + #[cfg(test)] + let permit_wait_pending = self.permit_wait_pending.clone(); + let (sender, mut receiver) = tokio::sync::mpsc::channel(1); 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 path = path.clone(); let stats = stats.clone(); + #[cfg(test)] + let permit_wait_pending = permit_wait_pending.clone(); async move { let permit_wait_start = stats.as_ref().map(|_| Instant::now()); - let permit = Arc::clone(&permits).acquire_owned().await; + let permit = Arc::clone(&permits).acquire_owned(); + tokio::pin!(permit); + #[cfg(test)] + let permit = if let Some(wait_pending) = permit_wait_pending { + let mut notified = false; + std::future::poll_fn(|context| { + let result = std::future::Future::poll(permit.as_mut(), context); + if result.is_pending() && !notified { + wait_pending.add_permits(1); + notified = true; + } + result + }) + .await + } else { + permit.await + }; + #[cfg(not(test))] + let permit = permit.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 permit = permit.map_err(|_| { io::Error::other("vindex range read concurrency limiter closed") })?; let expected = (range.end - range.start) as usize; @@ -282,24 +315,40 @@ impl VindexFileReader { .returned_bytes .fetch_add(data.len() as u64, Ordering::Relaxed); } - Ok(data) + Ok((index, data, permit)) } })) - .await; - let _ = sender.send(fetched); + .buffer_unordered(max_range_read_concurrency); + futures::pin_mut!(fetched); + while let Some(result) = fetched.next().await { + let failed = result.is_err(); + if sender.send(result).await.is_err() || failed { + 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, _permit) = fetched.ok_or_else(|| { + io::Error::other(format!( + "vindex range read task for '{}' was cancelled", + self.path + )) + })??; + consume(index, data) }); - 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<()> { @@ -358,8 +407,8 @@ impl VindexFileReader { .push(merged.len()); } let ranges: Vec<_> = merged.iter().map(|merged| merged.range.clone()).collect(); - let fetched = self.fetch_range_batch(&ranges)?; - for (merged_range, data) in merged.iter().zip(fetched) { + self.fetch_range_batch(&ranges, |merged_index, data| { + 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; @@ -367,8 +416,8 @@ impl VindexFileReader { .buf .copy_from_slice(&data[start..start + request.buf.len()]); } - } - Ok(()) + Ok(()) + }) } } @@ -413,6 +462,8 @@ impl SeekRead for VindexFileReader { path: self.path.clone(), scalar_cache: None, stats: self.stats.clone(), + #[cfg(test)] + permit_wait_pending: self.permit_wait_pending.clone(), })) } @@ -432,6 +483,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>>, @@ -478,6 +537,37 @@ mod tests { 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 OrderedRead { + data: Bytes, + ranges: Mutex>>, + started: tokio::sync::Semaphore, + release: tokio::sync::Semaphore, + } + #[async_trait] impl FileRead for RuntimeTrackingRead { async fn read(&self, range: Range) -> crate::Result { @@ -511,6 +601,41 @@ mod tests { } } + #[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 OrderedRead { + async fn read(&self, range: Range) -> crate::Result { + self.ranges.lock().unwrap().push(range.clone()); + self.started.add_permits(1); + acquire_test_permits( + &self.release, + 1, + "test did not release an ordered range read", + ) + .await; + Ok(self.data.slice(range.start as usize..range.end as usize)) + } + } + #[async_trait] impl FileRead for TrackingRead { async fn read(&self, range: Range) -> crate::Result { @@ -1002,6 +1127,119 @@ mod tests { assert_eq!(tracking.max_active.load(Ordering::SeqCst), concurrency); } + #[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_permits( + source, + tokio::runtime::Handle::current(), + Arc::new(tokio::sync::Semaphore::new(concurrency)), + 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 cloned_reader_is_not_queued_behind_an_entire_batch() { + let stride = RANGE_COALESCE_GAP + 2; + let data = Bytes::from(vec![8u8; 3 * stride as usize + 128]); + let tracking = Arc::new(OrderedRead { + data: data.clone(), + ranges: Mutex::new(Vec::new()), + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), + }); + let permits = Arc::new(tokio::sync::Semaphore::new(1)); + let source: Arc = tracking.clone(); + let mut first_reader = VindexFileReader::new_with_permits( + source, + tokio::runtime::Handle::current(), + permits, + 1, + data.len() as u64, + "index".to_string(), + ); + let mut second_reader = first_reader.try_clone_reader().unwrap().unwrap(); + let second_wait_pending = Arc::new(tokio::sync::Semaphore::new(0)); + second_reader.permit_wait_pending = Some(Arc::clone(&second_wait_pending)); + + let first = tokio::task::spawn_blocking(move || { + let mut first = [0u8; 1]; + let mut second = [0u8; 1]; + first_reader + .pread(&mut [ + ReadRequest::new(0, &mut first), + ReadRequest::new(stride, &mut second), + ]) + .unwrap(); + }); + acquire_test_permits(&tracking.started, 1, "first reader did not start").await; + + let second = tokio::task::spawn_blocking(move || { + let mut output = [0u8; 128]; + second_reader + .pread(&mut [ReadRequest::new(2 * stride, &mut output)]) + .unwrap(); + }); + acquire_test_permits( + &second_wait_pending, + 1, + "cloned reader did not queue for a range-read permit", + ) + .await; + + tracking.release.add_permits(1); + acquire_test_permits(&tracking.started, 1, "next range read did not start").await; + let second_started_range = tracking.ranges.lock().unwrap()[1].clone(); + tracking.release.add_permits(3); + tokio::time::timeout(Duration::from_secs(5), first) + .await + .expect("first reader task did not finish") + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), second) + .await + .expect("cloned reader task did not finish") + .unwrap(); + + assert_eq!(second_started_range.start, 2 * stride); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn failed_range_read_releases_permit() { let data = Bytes::from(vec![9u8; 1024]); diff --git a/docs/src/sql.md b/docs/src/sql.md index eaf9758a..58669e43 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.range-read-thread-num` | `32` | 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.range-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 `32`. + ### Variant Shredding Options Set these as table options when writing `VARIANT` columns to Parquet. The From 1e172b3b0916110ba94ca48638c7512798b081cb Mon Sep 17 00:00:00 2001 From: yantian Date: Mon, 17 Aug 2026 18:41:21 +0800 Subject: [PATCH 04/14] fix(vindex): track active range reads separately --- crates/paimon/src/vindex/range_reader.rs | 25 ++++++++++++++++-------- 1 file changed, 17 insertions(+), 8 deletions(-) diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 95341526..3729272e 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -58,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, @@ -67,6 +75,7 @@ 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, @@ -284,21 +293,21 @@ impl VindexFileReader { 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); - stats.peak_in_flight_reads.fetch_max( - max_range_read_concurrency - .saturating_sub(permits.available_permits()) - as u64, - Ordering::Relaxed, - ); - } + 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 data = reader.read(range.clone()).await.map_err(|error| { io::Error::other(format!( "failed to read vindex file '{path}' range {}..{}: {error}", range.start, range.end )) })?; + drop(in_flight_read); if data.len() != expected { return Err(io::Error::new( io::ErrorKind::UnexpectedEof, From 995f5fc3cc5737e442acb00700c4a4127efce969 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 10:56:57 +0800 Subject: [PATCH 05/14] perf(vindex): decouple range read I/O and response limits --- .../paimon/src/table/vector_search_builder.rs | 46 ++-- crates/paimon/src/vindex/range_reader.rs | 233 +++++++++++++----- ...8-17-vindex-range-read-double-buffering.md | 172 +++++++++++++ 3 files changed, 365 insertions(+), 86 deletions(-) create mode 100644 docs/plans/2026-08-17-vindex-range-read-double-buffering.md diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index ac3a813a..83f26fc9 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}; @@ -894,16 +894,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 = match backend { - VectorIndexBackend::Vindex => Some(Arc::new(tokio::sync::Semaphore::new( - range_read_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 = loader_range_read_permits.clone(); + let range_read_limiter = loader_range_read_limiter.clone(); let path = segment.path.clone(); let file_size = segment.file_size; Box::pin(async move { @@ -929,11 +927,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.expect("Vindex range-read permits"), - range_read_concurrency, + range_read_limiter.expect("Vindex range-read limiter"), file_size, path, ), @@ -1655,24 +1652,20 @@ async fn evaluate_batch_vector_search( .iter() .filter(|entry| is_vindex_index_type(&entry.index_file.index_type)) .count(); - let (batch_index_parallelism, range_read_concurrency, range_read_permits) = - if vindex_entry_count == 0 { - (1, 0, None) - } else { - let (index_parallelism, range_read_concurrency) = - vindex_concurrency_limits(&core_options, vindex_entry_count, concurrency)?; - ( - index_parallelism, - range_read_concurrency, - Some(Arc::new(tokio::sync::Semaphore::new( - range_read_concurrency, - ))), - ) - }; + 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 = range_read_permits.clone(); + 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"); @@ -1753,11 +1746,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.expect("Vindex range-read permits"), - range_read_concurrency, + range_read_limiter.expect("Vindex range-read limiter"), file_size, file_name.clone(), ); diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 3729272e..87365193 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -115,14 +115,35 @@ impl RangeIoStats { } } +#[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, + } + } +} + /// 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, - max_range_read_concurrency: usize, + limiter: RangeReadLimiter, file_size: u64, path: String, scalar_cache: Option, @@ -139,31 +160,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( - DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM, - )), - DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM, + RangeReadLimiter::new(DEFAULT_GLOBAL_INDEX_RANGE_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, - max_range_read_concurrency: usize, + limiter: RangeReadLimiter, file_size: u64, path: String, ) -> Self { Self { reader, runtime, - permits, - max_range_read_concurrency, + limiter, file_size, path, scalar_cache: None, @@ -247,25 +263,35 @@ impl VindexFileReader { mut consume: impl FnMut(usize, Bytes) -> 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 max_range_read_concurrency = self.max_range_read_concurrency; + let io_limit = self.limiter.io_limit; + let response_limit = self.limiter.response_limit; #[cfg(test)] let permit_wait_pending = self.permit_wait_pending.clone(); - let (sender, mut receiver) = tokio::sync::mpsc::channel(1); + let (sender, mut receiver) = tokio::sync::mpsc::channel(io_limit); self.runtime.spawn(async move { 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(); #[cfg(test)] let permit_wait_pending = permit_wait_pending.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 = Arc::clone(&permits).acquire_owned(); + let permit = io_permits.acquire_owned(); tokio::pin!(permit); #[cfg(test)] let permit = if let Some(wait_pending) = permit_wait_pending { @@ -289,7 +315,7 @@ impl VindexFileReader { .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; @@ -301,13 +327,15 @@ impl VindexFileReader { .fetch_max(active, Ordering::Relaxed); InFlightRead(&stats.in_flight_reads) }); - let data = reader.read(range.clone()).await.map_err(|error| { + 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 )) })?; - drop(in_flight_read); if data.len() != expected { return Err(io::Error::new( io::ErrorKind::UnexpectedEof, @@ -324,14 +352,17 @@ impl VindexFileReader { .returned_bytes .fetch_add(data.len() as u64, Ordering::Relaxed); } - Ok((index, data, permit)) + sender + .send(Ok((index, data, response_permit))) + .await + .map_err(|_| io::Error::other("vindex range read receiver closed")) } })) - .buffer_unordered(max_range_read_concurrency); + .buffer_unordered(response_limit); futures::pin_mut!(fetched); while let Some(result) = fetched.next().await { - let failed = result.is_err(); - if sender.send(result).await.is_err() || failed { + if let Err(error) = result { + let _ = sender.send(Err(error)).await; return; } } @@ -344,7 +375,7 @@ impl VindexFileReader { if let Some(start) = wait_start { io_wait_nanos = io_wait_nanos.saturating_add(start.elapsed().as_nanos() as u64); } - let (index, data, _permit) = fetched.ok_or_else(|| { + let (index, data, _response_permit) = fetched.ok_or_else(|| { io::Error::other(format!( "vindex range read task for '{}' was cancelled", self.path @@ -465,8 +496,7 @@ impl SeekRead for VindexFileReader { Ok(Some(Self { reader: Arc::clone(&self.reader), runtime: self.runtime.clone(), - permits: Arc::clone(&self.permits), - max_range_read_concurrency: self.max_range_read_concurrency, + limiter: self.limiter.clone(), file_size: self.file_size, path: self.path.clone(), scalar_cache: None, @@ -971,20 +1001,26 @@ mod tests { started: tokio::sync::Semaphore::new(0), release: tokio::sync::Semaphore::new(0), }); - let permits = Arc::new(tokio::sync::Semaphore::new(1)); + let limiter = RangeReadLimiter::new(1); let source: Arc = tracking.clone(); - let mut first_reader = VindexFileReader::new_with_permits( + let mut first_reader = VindexFileReader::new_with_limiter( source, tokio::runtime::Handle::current(), - Arc::clone(&permits), - 1, + limiter, data.len() as u64, "index".to_string(), ); let stats = Arc::new(RangeIoStats::default()); first_reader.stats = Some(Arc::clone(&stats)); let mut second_reader = first_reader.try_clone_reader().unwrap().unwrap(); - assert!(Arc::ptr_eq(&first_reader.permits, &second_reader.permits)); + 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]; @@ -1013,6 +1049,20 @@ 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 configured_range_read_concurrency_can_exceed_32() { let configured_concurrency = 64; @@ -1026,13 +1076,11 @@ mod tests { started: tokio::sync::Semaphore::new(0), release: tokio::sync::Semaphore::new(0), }); - let permits = Arc::new(tokio::sync::Semaphore::new(configured_concurrency)); let source: Arc = tracking.clone(); - let mut reader = VindexFileReader::new_with_permits( + let mut reader = VindexFileReader::new_with_limiter( source, tokio::runtime::Handle::current(), - permits, - configured_concurrency, + RangeReadLimiter::new(configured_concurrency), data.len() as u64, "index".to_string(), ); @@ -1097,13 +1145,11 @@ mod tests { started: tokio::sync::Semaphore::new(0), release: tokio::sync::Semaphore::new(0), }); - let permits = Arc::new(tokio::sync::Semaphore::new(concurrency)); let source: Arc = tracking.clone(); - let mut reader = VindexFileReader::new_with_permits( + let mut reader = VindexFileReader::new_with_limiter( source, tokio::runtime::Handle::current(), - permits, - concurrency, + RangeReadLimiter::new(concurrency), data.len() as u64, "index".to_string(), ); @@ -1136,6 +1182,72 @@ mod tests { 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 || { + second_reader + .fetch_range_batch(&[2 * stride..2 * stride + 1], |_, _| 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 completed_range_buffers_are_released_before_slowest_read() { let concurrency = 2; @@ -1149,11 +1261,10 @@ mod tests { dropped: Arc::clone(&dropped), }); let source: Arc = tracking.clone(); - let mut reader = VindexFileReader::new_with_permits( + let mut reader = VindexFileReader::new_with_limiter( source, tokio::runtime::Handle::current(), - Arc::new(tokio::sync::Semaphore::new(concurrency)), - concurrency, + RangeReadLimiter::new(concurrency), range_count as u64 * stride, "index".to_string(), ); @@ -1187,20 +1298,18 @@ mod tests { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn cloned_reader_is_not_queued_behind_an_entire_batch() { let stride = RANGE_COALESCE_GAP + 2; - let data = Bytes::from(vec![8u8; 3 * stride as usize + 128]); + let data = Bytes::from(vec![8u8; 4 * stride as usize + 128]); let tracking = Arc::new(OrderedRead { data: data.clone(), ranges: Mutex::new(Vec::new()), started: tokio::sync::Semaphore::new(0), release: tokio::sync::Semaphore::new(0), }); - let permits = Arc::new(tokio::sync::Semaphore::new(1)); let source: Arc = tracking.clone(); - let mut first_reader = VindexFileReader::new_with_permits( + let mut first_reader = VindexFileReader::new_with_limiter( source, tokio::runtime::Handle::current(), - permits, - 1, + RangeReadLimiter::new(1), data.len() as u64, "index".to_string(), ); @@ -1211,10 +1320,12 @@ mod tests { let first = tokio::task::spawn_blocking(move || { let mut first = [0u8; 1]; let mut second = [0u8; 1]; + let mut third = [0u8; 1]; first_reader .pread(&mut [ ReadRequest::new(0, &mut first), ReadRequest::new(stride, &mut second), + ReadRequest::new(2 * stride, &mut third), ]) .unwrap(); }); @@ -1223,19 +1334,21 @@ mod tests { let second = tokio::task::spawn_blocking(move || { let mut output = [0u8; 128]; second_reader - .pread(&mut [ReadRequest::new(2 * stride, &mut output)]) + .pread(&mut [ReadRequest::new(3 * stride, &mut output)]) .unwrap(); }); + + tracking.release.add_permits(1); + acquire_test_permits(&tracking.started, 1, "second range read did not start").await; acquire_test_permits( &second_wait_pending, 1, - "cloned reader did not queue for a range-read permit", + "cloned reader did not queue within the response window", ) .await; - tracking.release.add_permits(1); - acquire_test_permits(&tracking.started, 1, "next range read did not start").await; - let second_started_range = tracking.ranges.lock().unwrap()[1].clone(); + acquire_test_permits(&tracking.started, 1, "cloned range read did not start").await; + let third_started_range = tracking.ranges.lock().unwrap()[2].clone(); tracking.release.add_permits(3); tokio::time::timeout(Duration::from_secs(5), first) .await @@ -1246,7 +1359,7 @@ mod tests { .expect("cloned reader task did not finish") .unwrap(); - assert_eq!(second_started_range.start, 2 * stride); + assert_eq!(third_started_range.start, 3 * stride); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -1256,12 +1369,13 @@ mod tests { data: data.clone(), calls: AtomicUsize::new(0), }); - let permits = Arc::new(tokio::sync::Semaphore::new(1)); - let mut reader = VindexFileReader::new_with_permits( + 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(), - Arc::clone(&permits), - 1, + limiter, data.len() as u64, "index".to_string(), ); @@ -1276,7 +1390,8 @@ mod tests { .await .unwrap(); assert_eq!(error.kind(), io::ErrorKind::Other); - assert_eq!(permits.available_permits(), 1); + assert_eq!(io_permits.available_permits(), 1); + assert_eq!(response_permits.available_permits(), 2); tokio::time::timeout( Duration::from_secs(5), diff --git a/docs/plans/2026-08-17-vindex-range-read-double-buffering.md b/docs/plans/2026-08-17-vindex-range-read-double-buffering.md new file mode 100644 index 00000000..d65b54a9 --- /dev/null +++ b/docs/plans/2026-08-17-vindex-range-read-double-buffering.md @@ -0,0 +1,172 @@ +# Vindex range-read 双限流优化计划 + +> 状态(2026-08-18):已在 `perf/vindex-range-read-concurrency` 工作树实施,尚未提交。本文件现作为实施与验收记录;基线 commit 仍为 `1e172b3`,下列 Task 6 A/B 尚待执行。 + +## 1. 结论 + +当前 `1e172b3` 已把响应内存从“保留全部 range”降为有界,但同一个 semaphore 同时覆盖底层 I/O、channel 排队和请求 buffer 复制。热 cache 下读取很快,permit 的主要持有时间会转移到交付和复制阶段,因此真实 I/O 不能及时补位;这与 4 GiB 热 cache 测试中约 3.4% 的 QPS 回退方向一致。 + +采用两个独立限制是正确方向: + +- I/O limit:`C = global-index.range-read-thread-num`,只覆盖 `FileRead::read(...).await`。 +- Response limit:内部固定为 `min(2C, Semaphore::MAX_PERMITS)`,覆盖从发起读取前到响应被消费完成的生命周期。 +- Channel capacity:`C`,允许已完成响应批量交给同步消费者。 + +但不能只把第二个 semaphore 和 channel 加到当前实现。当前 producer 在 `buffer_unordered(C)` 外逐个 `sender.send(...).await`;channel 满时 producer 停止轮询整个 stream,其他已就绪 read future 也不能及时完成和释放 I/O permit。应把 `send` 放进每个受限 work item,并让 work window 使用 response limit(`2C`),再由 I/O semaphore 单独限制真实读取为 `C`。 + +### 1.1 新增 benchmark 证据(2026-08-17) + +| 实现 | Warmup | 50q 耗时 | QPS | Recall | NDCG | +|---|---:|---:|---:|---:|---:| +| 旧版 `0946fcd` | 3.669s | 3.105s | 16.101 | 0.8182 | 0.86104 | +| 标记为新版 `ab90d01` | 3.810s | 3.761s | 13.295 | 0.8182 | 0.86104 | + +该样本中新版 Warmup 慢 3.8%,50q 耗时增加 21.1%,QPS 下降 17.4%,已经超过单轮约 3% 的噪声容忍范围,必须作为性能回退候选处理;Recall/NDCG 一致,说明不是搜索参数或结果质量变化换来的差异。 + +不过 `ab90d01` 是当前 PR 的 main 基线,不是当前 feature HEAD `1e172b3`。正式归因前必须从编译日志或 binary provenance 确认实际源码 commit,不能只相信结果 JSON/表格中的硬编码字段。Lance 的 Recall/NDCG 与 Paimon 不同,只作为外部吞吐参考,不用于计算本 PR 的回退幅度。 + +## 2. 必须保持的边界 + +1. 两个 semaphore 都必须在一次 vector search 的 reader 创建点生成,并被所有 index reader 和 clone 共享;不得按 file、clone 或 `pread` 单独创建,否则总响应内存会随 reader 数放大。 +2. 获取顺序固定为 response permit -> I/O permit。反过来会让任务持有稀缺 I/O permit 等待内存预算。 +3. `FileRead::read` 返回后立即释放 I/O permit;长度校验、channel 等待和请求 buffer 复制不得占用它。 +4. Response permit 随 `(index, Bytes)` 穿过 channel,在 `consume` 返回后释放;错误、receiver 关闭和 future 取消也必须依靠 RAII 释放。 +5. `range_permit_wait_nanos` 继续只统计 I/O permit 等待,不混入 response backpressure;`peak_in_flight_reads` 继续由 `InFlightRead` 统计真实底层读取。 +6. 该限制按响应个数而不是字节数计量。瞬时 multi-range 响应为 `O(2C * max_response_size)`;请求方 buffer 和既有的每-reader 64 KiB scalar read-ahead cache 不在此预算内。 +7. `2C` 是吞吐/内存折中,不保证任何时刻严格跑满 `C` 路 I/O。消费者、channel 和等待发送的响应耗尽窗口时会自然背压;不为理论上的单槽气泡引入 `2C+1` 或第三个配置,除非 benchmark 证明有必要。 + +## 3. 修改范围 + +只修改: + +- `crates/paimon/src/vindex/range_reader.rs` +- `crates/paimon/src/table/vector_search_builder.rs` + +不修改配置、SQL 文档、vindex-core、合并规则或现有诊断输出格式;不增加依赖和用户参数。 + +## 4. 实施步骤 + +### Task 0:锁定基线(已完成) + +实施基线为分支 `perf/vindex-range-read-concurrency` 的 `1e172b3`。当前工作树已有本计划的未提交实现,因此不再预期干净状态。 + +```bash +git status --short --branch +git rev-parse HEAD +``` + +实施前确认 HEAD 为 `1e172b3`;当前两个目标文件的修改即本计划候选实现。 + +### Task 1:行为测试(已完成,实施后对账) + +文件:`crates/paimon/src/vindex/range_reader.rs` 的现有 tests 模块。 + +1. `response_copy_and_buffers_are_bounded_across_clones` 使用 `C=1`,让第一个 reader 的响应进入 `consume` 后阻塞,验证第二个底层 read 已启动,同时另一个 clone 的第 `2C+1` 个 read 在释放 response permit 前不会启动、释放后能够继续。 +2. 该测试合并覆盖“复制不持有 I/O permit”和“所有 clone 全局共享 2C 响应上限”,避免两套重复测试样板。 +3. 正向同步使用 semaphore/channel;“第 `2C+1` 个 read 尚未启动”是必要的否定断言,明确允许使用 1 秒 timeout 作为例外。其余 timeout 只用于防止 CI 永久挂起,不作为性能阈值。 + +本计划与已有工作树对账时生产实现已经存在,无法在当前状态重放 RED-first;这里记录为实施后的回归/characterization test,不伪造 RED 结果。 + +验证: + +```bash +cargo test -p paimon vindex::range_reader::tests::response_copy_and_buffers_are_bounded_across_clones +``` + +### Task 2:把两个限制绑定为一个共享内部值(已完成) + +文件:`crates/paimon/src/vindex/range_reader.rs`。 + +1. 用一个最小的 `pub(crate) RangeReadLimiter` 绑定 shared I/O semaphore、shared response semaphore、`io_limit=C` 和 `response_limit=min(C.saturating_mul(2), Semaphore::MAX_PERMITS)`。 +2. `VindexFileReader` 持有 limiter,替换当前分离的 `permits + max_range_read_concurrency`,避免构造时传入互相矛盾的 semaphore size 和并发值。 +3. 测试构造器仍用默认 `C=32`;`try_clone_reader` clone 同一个 limiter 内的两个 `Arc`。 +4. `response_limit_saturates_at_semaphore_max_permits` 覆盖 `C > MAX_PERMITS / 2` 时不会乘法溢出或让 `Semaphore::new` panic。 + +验证:现有 `clones_share_range_read_permits` 更新为同时断言 I/O 与 response semaphore 被共享;不增加通用 limiter trait、factory 或新配置。 + +### Task 3:在两个生产入口全局共享 limiter(已完成) + +文件:`crates/paimon/src/table/vector_search_builder.rs`。 + +在以下两个现有 semaphore 创建点各创建一次 `RangeReadLimiter::new(range_read_concurrency)`,再 clone 给所有 Vindex reader: + +- `plan_and_search_pk_candidates_batch` 的 loader 路径; +- `evaluate_batch_vector_search` 的 vector entry 路径。 + +Lumina 路径保持不变。`global-index.thread-num`、native batch parallelism 和 `global-index.range-read-thread-num` 的现有语义保持不变。 + +验证:`rg -n "RangeReadLimiter::new" crates/paimon/src/table/vector_search_builder.rs` 只应命中上述两个 search scope,不能出现在 per-file closure 内。 + +### Task 4:重排 fetch pipeline(已完成) + +文件:`crates/paimon/src/vindex/range_reader.rs` 的 `fetch_range_batch`。 + +每个 range work item 按以下顺序执行: + +```text +acquire response permit + -> acquire I/O permit + -> FileRead::read + -> drop InFlightRead and I/O permit + -> validate response + -> bounded channel send + -> synchronous consume/copy + -> drop response permit +``` + +具体要求: + +1. channel capacity 从 `1` 改为 `C`。 +2. 把 `sender.send` 移进每个 range future;future 完成发送后才算 work item 完成。 +3. `buffer_unordered` 的 work window 使用 `response_limit`,真实读取仍由 I/O semaphore 限制为 `C`。 +4. 第一个 read/send 错误送达 consumer 后停止 producer;drop 剩余 stream,取消未完成 future 并释放两个 permit。 +5. `FileRead::read(...).await` 的结果先保存,再立即 drop `InFlightRead` 和 I/O permit,然后执行错误映射、长度校验和发送。 +6. 不改变 completion-order copy、range merge、scalar cache 或 stats 字段。 + +### Task 5:补齐回归测试(已完成) + +更新并保留以下已有测试语义: + +- `clones_share_range_read_permits`:两个限制均跨 clone 共享。 +- `configured_range_read_concurrency_can_exceed_32`:真实 peak I/O 仍能达到配置值。 +- `range_reads_refill_before_slowest_batch_member_finishes`:释放 I/O 后及时补位。 +- `completed_range_buffers_are_released_before_slowest_read`:消费完成的 payload 及时 drop。 +- `cloned_reader_is_not_queued_behind_an_entire_batch`:多个 reader 仍共享且公平竞争 I/O。 +- `failed_range_read_releases_permit`:改为同时验证 I/O 与 response permit,后续读取不死锁。 +- `range_io_stats_count_coalesced_reads`:`peak_in_flight_reads` 仍统计真实 read,而不是 response 生命周期。 + +运行: + +```bash +cargo fmt --all --check +cargo test -p paimon vindex::range_reader +cargo clippy -p paimon --all-targets -- -D warnings +git diff --check +``` + +本机 Apple ARM 若仍被 `paimon-vindex-core 0.3.0` 的 unstable NEON 编译问题阻塞,至少本地完成 `fmt` 和 `diff --check`,完整 test/clippy 以 Linux GitHub CI 为门禁;不得把依赖绕过改动混入本 PR。 + +2026-08-18 本地验证记录(Rust `1.97.0`): + +- `cargo fmt --all -- --check` 通过; +- `cargo test -p paimon vindex::range_reader::tests --lib`:20 passed; +- `cargo clippy -p paimon --all-targets -- -D warnings` 通过; +- `git diff --check` 通过。 + +### Task 6:复跑现有 A/B(待执行) + +用同一 benchmark 配置分别比较基线 `1e172b3` 与候选提交,每组至少 3 次取中位数: + +1. 无 local cache、无 warmup:确认冷读 QPS/平均延迟变化在约 3% 测试噪声内,Recall/NDCG 完全一致。 +2. 4 GiB memory cache、热读:候选版本至少不能比 `1e172b3` 更慢;目标是收回当前相对 `0946fcd` 约 3.4% 的 QPS 差距。 +3. 记录 peak RSS、`io_wait`、`range_permit_wait` 和 `peak_in_flight_reads`。RSS 可以高于严格 `C` 响应上限的 `1e172b3`,但不能回到保留全部 range 的增长模式;全局 `2C` 上限由确定性单测负责证明。 + +如果热 cache 仍稳定回退超过 3%,先保留双 limiter 的正确性改动但不宣称性能改善,使用现有 timing 数据确认瓶颈后再决定是否调整窗口;本计划不预先增加自适应窗口或 byte-based limiter。 + +## 5. 完成标准 + +- 真实底层 read 的 peak 不超过 `C`,响应 payload 数全局不超过 `min(2C, MAX_PERMITS)`。 +- 同步复制不再持有 I/O permit,clone/file 之间共享两种预算。 +- 错误、取消和 receiver 关闭不会泄漏 permit 或挂住后续读取。 +- Recall/NDCG 不变,冷读无稳定回退,热 cache 回退被消除或有 timing 证据解释。 +- PR 只新增上述两文件的实现/测试改动,不引入新配置、依赖或相邻重构。 From 836af2cbc7d3fe29dce392ed3955cd2661e1d153 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 11:25:25 +0800 Subject: [PATCH 06/14] fix(ci): stabilize range reader concurrency test --- crates/paimon/src/vindex/range_reader.rs | 57 ++++++++----------- ...8-17-vindex-range-read-double-buffering.md | 21 ++++++- 2 files changed, 43 insertions(+), 35 deletions(-) diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 87365193..401ff005 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -149,7 +149,7 @@ pub(crate) struct VindexFileReader { scalar_cache: Option, stats: Option>, #[cfg(test)] - permit_wait_pending: Option>, + response_permit_wait_started: Option>, } impl VindexFileReader { @@ -185,7 +185,7 @@ impl VindexFileReader { scalar_cache: None, stats: vector_search_timing_enabled().then(|| Arc::new(RangeIoStats::default())), #[cfg(test)] - permit_wait_pending: None, + response_permit_wait_started: None, } } @@ -271,7 +271,7 @@ impl VindexFileReader { let io_limit = self.limiter.io_limit; let response_limit = self.limiter.response_limit; #[cfg(test)] - let permit_wait_pending = self.permit_wait_pending.clone(); + let response_permit_wait_started = self.response_permit_wait_started.clone(); let (sender, mut receiver) = tokio::sync::mpsc::channel(io_limit); self.runtime.spawn(async move { let fetched = stream::iter(requested.into_iter().enumerate().map(|(index, range)| { @@ -282,8 +282,12 @@ impl VindexFileReader { let path = path.clone(); let stats = stats.clone(); #[cfg(test)] - let permit_wait_pending = permit_wait_pending.clone(); + let response_permit_wait_started = response_permit_wait_started.clone(); async move { + #[cfg(test)] + if let Some(wait_started) = response_permit_wait_started { + wait_started.add_permits(1); + } let response_permit = response_permits .acquire_owned() .await @@ -291,25 +295,7 @@ impl VindexFileReader { io::Error::other("vindex range response limiter closed") })?; let permit_wait_start = stats.as_ref().map(|_| Instant::now()); - let permit = io_permits.acquire_owned(); - tokio::pin!(permit); - #[cfg(test)] - let permit = if let Some(wait_pending) = permit_wait_pending { - let mut notified = false; - std::future::poll_fn(|context| { - let result = std::future::Future::poll(permit.as_mut(), context); - if result.is_pending() && !notified { - wait_pending.add_permits(1); - notified = true; - } - result - }) - .await - } else { - permit.await - }; - #[cfg(not(test))] - let permit = permit.await; + let permit = io_permits.acquire_owned().await; if let (Some(stats), Some(start)) = (&stats, permit_wait_start) { stats .range_permit_wait_nanos @@ -502,7 +488,7 @@ impl SeekRead for VindexFileReader { scalar_cache: None, stats: self.stats.clone(), #[cfg(test)] - permit_wait_pending: self.permit_wait_pending.clone(), + response_permit_wait_started: self.response_permit_wait_started.clone(), })) } @@ -1224,8 +1210,9 @@ mod tests { .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(&[2 * stride..2 * stride + 1], |_, _| Ok(())) + .fetch_range_batch(std::slice::from_ref(&range), |_, _| Ok(())) .unwrap(); }); tracking.release.add_permits(1); @@ -1314,8 +1301,8 @@ mod tests { "index".to_string(), ); let mut second_reader = first_reader.try_clone_reader().unwrap().unwrap(); - let second_wait_pending = Arc::new(tokio::sync::Semaphore::new(0)); - second_reader.permit_wait_pending = Some(Arc::clone(&second_wait_pending)); + let response_wait_started = Arc::new(tokio::sync::Semaphore::new(0)); + second_reader.response_permit_wait_started = Some(Arc::clone(&response_wait_started)); let first = tokio::task::spawn_blocking(move || { let mut first = [0u8; 1]; @@ -1337,18 +1324,18 @@ mod tests { .pread(&mut [ReadRequest::new(3 * stride, &mut output)]) .unwrap(); }); - - tracking.release.add_permits(1); - acquire_test_permits(&tracking.started, 1, "second range read did not start").await; acquire_test_permits( - &second_wait_pending, + &response_wait_started, 1, - "cloned reader did not queue within the response window", + "cloned reader did not start waiting for a response permit", ) .await; + + tracking.release.add_permits(1); + acquire_test_permits(&tracking.started, 1, "second range read did not start").await; tracking.release.add_permits(1); acquire_test_permits(&tracking.started, 1, "cloned range read did not start").await; - let third_started_range = tracking.ranges.lock().unwrap()[2].clone(); + let first_three_ranges = tracking.ranges.lock().unwrap()[..3].to_vec(); tracking.release.add_permits(3); tokio::time::timeout(Duration::from_secs(5), first) .await @@ -1359,7 +1346,9 @@ mod tests { .expect("cloned reader task did not finish") .unwrap(); - assert_eq!(third_started_range.start, 3 * stride); + assert!(first_three_ranges + .iter() + .any(|range| range.start == 3 * stride)); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] diff --git a/docs/plans/2026-08-17-vindex-range-read-double-buffering.md b/docs/plans/2026-08-17-vindex-range-read-double-buffering.md index d65b54a9..0387225e 100644 --- a/docs/plans/2026-08-17-vindex-range-read-double-buffering.md +++ b/docs/plans/2026-08-17-vindex-range-read-double-buffering.md @@ -1,6 +1,25 @@ + + # Vindex range-read 双限流优化计划 -> 状态(2026-08-18):已在 `perf/vindex-range-read-concurrency` 工作树实施,尚未提交。本文件现作为实施与验收记录;基线 commit 仍为 `1e172b3`,下列 Task 6 A/B 尚待执行。 +> 状态(2026-08-18):已在 `perf/vindex-range-read-concurrency` 实施。本文件现作为实施与验收记录;基线 commit 仍为 `1e172b3`,下列 Task 6 A/B 尚待执行。 ## 1. 结论 From 819c7a9e2eeb73a0ce140a66a5b356bba5407dbd Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 11:42:31 +0800 Subject: [PATCH 07/14] fix(vindex): bound range read stats memory --- .../paimon/src/table/vector_search_builder.rs | 6 +- crates/paimon/src/vindex/range_reader.rs | 63 ++++++++++++++----- 2 files changed, 51 insertions(+), 18 deletions(-) diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 83f26fc9..b1d023b1 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -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} peak_in_flight_reads={} read_many_merged_ranges={} read_many_chunks={} read_many_chunk_sizes={:?}", + "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, @@ -157,7 +157,9 @@ fn log_vindex_range_io_stats(file: &str, query_count: usize, stats: &RangeIoStat stats.peak_in_flight_reads, stats.read_many_merged_ranges, stats.read_many_chunks, - stats.read_many_chunk_sizes, + stats.read_many_chunk_size_sum, + stats.read_many_chunk_size_min, + stats.read_many_chunk_size_max, ); } diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 401ff005..edfdd1e9 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -79,7 +79,9 @@ pub(crate) struct RangeIoStats { peak_in_flight_reads: AtomicU64, read_many_merged_ranges: AtomicU64, read_many_chunks: AtomicU64, - read_many_chunk_sizes: std::sync::Mutex>, + read_many_chunk_size_sum: AtomicU64, + read_many_chunk_size_min: AtomicU64, + read_many_chunk_size_max: AtomicU64, } #[derive(Clone, Debug, Default, PartialEq, Eq)] @@ -94,7 +96,9 @@ pub(crate) struct RangeIoStatsSnapshot { 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_sizes: Vec, + 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 { @@ -110,7 +114,9 @@ impl RangeIoStats { 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_sizes: self.read_many_chunk_sizes.lock().unwrap().clone(), + 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), } } } @@ -420,17 +426,28 @@ impl VindexFileReader { } if let Some(stats) = &self.stats { + let chunk_size = merged.len() as u64; stats .read_many_merged_ranges - .fetch_add(merged.len() as u64, Ordering::Relaxed); - } - if let Some(stats) = &self.stats { + .fetch_add(chunk_size, Ordering::Relaxed); stats.read_many_chunks.fetch_add(1, Ordering::Relaxed); stats - .read_many_chunk_sizes - .lock() - .unwrap() - .push(merged.len()); + .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); } let ranges: Vec<_> = merged.iter().map(|merged| merged.range.clone()).collect(); self.fetch_range_batch(&ranges, |merged_index, data| { @@ -924,9 +941,11 @@ mod tests { stats.peak_in_flight_reads, stats.read_many_merged_ranges, stats.read_many_chunks, - stats.read_many_chunk_sizes, + stats.read_many_chunk_size_sum, + stats.read_many_chunk_size_min, + stats.read_many_chunk_size_max, ), - (2, 256, 2, 256, 0, 1, 0, 0, Vec::new()) + (2, 256, 2, 256, 0, 1, 0, 0, 0, 0, 0) ); assert!(stats.io_wait_nanos > 0); } @@ -955,6 +974,14 @@ 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(); @@ -970,9 +997,11 @@ mod tests { stats.peak_in_flight_reads, stats.read_many_merged_ranges, stats.read_many_chunks, - stats.read_many_chunk_sizes, + stats.read_many_chunk_size_sum, + stats.read_many_chunk_size_min, + stats.read_many_chunk_size_max, ), - (3, 12, 2, 16, 0, 1, 2, 1, vec![2]) + (5, 20, 3, 28, 0, 1, 3, 2, 3, 1, 2) ); assert!(stats.io_wait_nanos > 0); } @@ -1110,9 +1139,11 @@ mod tests { snapshot.peak_in_flight_reads, snapshot.read_many_merged_ranges, snapshot.read_many_chunks, - snapshot.read_many_chunk_sizes, + 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, vec![65]) + (65, 65, 65, 65, 0, 64, 65, 1, 65, 65, 65) ); assert!(snapshot.io_wait_nanos > 0); assert!(snapshot.range_permit_wait_nanos > 0); From 4a9c484754ba527b87fc51858a96035623908c97 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 11:47:41 +0800 Subject: [PATCH 08/14] fix --- ...8-17-vindex-range-read-double-buffering.md | 191 ------------------ 1 file changed, 191 deletions(-) delete mode 100644 docs/plans/2026-08-17-vindex-range-read-double-buffering.md diff --git a/docs/plans/2026-08-17-vindex-range-read-double-buffering.md b/docs/plans/2026-08-17-vindex-range-read-double-buffering.md deleted file mode 100644 index 0387225e..00000000 --- a/docs/plans/2026-08-17-vindex-range-read-double-buffering.md +++ /dev/null @@ -1,191 +0,0 @@ - - -# Vindex range-read 双限流优化计划 - -> 状态(2026-08-18):已在 `perf/vindex-range-read-concurrency` 实施。本文件现作为实施与验收记录;基线 commit 仍为 `1e172b3`,下列 Task 6 A/B 尚待执行。 - -## 1. 结论 - -当前 `1e172b3` 已把响应内存从“保留全部 range”降为有界,但同一个 semaphore 同时覆盖底层 I/O、channel 排队和请求 buffer 复制。热 cache 下读取很快,permit 的主要持有时间会转移到交付和复制阶段,因此真实 I/O 不能及时补位;这与 4 GiB 热 cache 测试中约 3.4% 的 QPS 回退方向一致。 - -采用两个独立限制是正确方向: - -- I/O limit:`C = global-index.range-read-thread-num`,只覆盖 `FileRead::read(...).await`。 -- Response limit:内部固定为 `min(2C, Semaphore::MAX_PERMITS)`,覆盖从发起读取前到响应被消费完成的生命周期。 -- Channel capacity:`C`,允许已完成响应批量交给同步消费者。 - -但不能只把第二个 semaphore 和 channel 加到当前实现。当前 producer 在 `buffer_unordered(C)` 外逐个 `sender.send(...).await`;channel 满时 producer 停止轮询整个 stream,其他已就绪 read future 也不能及时完成和释放 I/O permit。应把 `send` 放进每个受限 work item,并让 work window 使用 response limit(`2C`),再由 I/O semaphore 单独限制真实读取为 `C`。 - -### 1.1 新增 benchmark 证据(2026-08-17) - -| 实现 | Warmup | 50q 耗时 | QPS | Recall | NDCG | -|---|---:|---:|---:|---:|---:| -| 旧版 `0946fcd` | 3.669s | 3.105s | 16.101 | 0.8182 | 0.86104 | -| 标记为新版 `ab90d01` | 3.810s | 3.761s | 13.295 | 0.8182 | 0.86104 | - -该样本中新版 Warmup 慢 3.8%,50q 耗时增加 21.1%,QPS 下降 17.4%,已经超过单轮约 3% 的噪声容忍范围,必须作为性能回退候选处理;Recall/NDCG 一致,说明不是搜索参数或结果质量变化换来的差异。 - -不过 `ab90d01` 是当前 PR 的 main 基线,不是当前 feature HEAD `1e172b3`。正式归因前必须从编译日志或 binary provenance 确认实际源码 commit,不能只相信结果 JSON/表格中的硬编码字段。Lance 的 Recall/NDCG 与 Paimon 不同,只作为外部吞吐参考,不用于计算本 PR 的回退幅度。 - -## 2. 必须保持的边界 - -1. 两个 semaphore 都必须在一次 vector search 的 reader 创建点生成,并被所有 index reader 和 clone 共享;不得按 file、clone 或 `pread` 单独创建,否则总响应内存会随 reader 数放大。 -2. 获取顺序固定为 response permit -> I/O permit。反过来会让任务持有稀缺 I/O permit 等待内存预算。 -3. `FileRead::read` 返回后立即释放 I/O permit;长度校验、channel 等待和请求 buffer 复制不得占用它。 -4. Response permit 随 `(index, Bytes)` 穿过 channel,在 `consume` 返回后释放;错误、receiver 关闭和 future 取消也必须依靠 RAII 释放。 -5. `range_permit_wait_nanos` 继续只统计 I/O permit 等待,不混入 response backpressure;`peak_in_flight_reads` 继续由 `InFlightRead` 统计真实底层读取。 -6. 该限制按响应个数而不是字节数计量。瞬时 multi-range 响应为 `O(2C * max_response_size)`;请求方 buffer 和既有的每-reader 64 KiB scalar read-ahead cache 不在此预算内。 -7. `2C` 是吞吐/内存折中,不保证任何时刻严格跑满 `C` 路 I/O。消费者、channel 和等待发送的响应耗尽窗口时会自然背压;不为理论上的单槽气泡引入 `2C+1` 或第三个配置,除非 benchmark 证明有必要。 - -## 3. 修改范围 - -只修改: - -- `crates/paimon/src/vindex/range_reader.rs` -- `crates/paimon/src/table/vector_search_builder.rs` - -不修改配置、SQL 文档、vindex-core、合并规则或现有诊断输出格式;不增加依赖和用户参数。 - -## 4. 实施步骤 - -### Task 0:锁定基线(已完成) - -实施基线为分支 `perf/vindex-range-read-concurrency` 的 `1e172b3`。当前工作树已有本计划的未提交实现,因此不再预期干净状态。 - -```bash -git status --short --branch -git rev-parse HEAD -``` - -实施前确认 HEAD 为 `1e172b3`;当前两个目标文件的修改即本计划候选实现。 - -### Task 1:行为测试(已完成,实施后对账) - -文件:`crates/paimon/src/vindex/range_reader.rs` 的现有 tests 模块。 - -1. `response_copy_and_buffers_are_bounded_across_clones` 使用 `C=1`,让第一个 reader 的响应进入 `consume` 后阻塞,验证第二个底层 read 已启动,同时另一个 clone 的第 `2C+1` 个 read 在释放 response permit 前不会启动、释放后能够继续。 -2. 该测试合并覆盖“复制不持有 I/O permit”和“所有 clone 全局共享 2C 响应上限”,避免两套重复测试样板。 -3. 正向同步使用 semaphore/channel;“第 `2C+1` 个 read 尚未启动”是必要的否定断言,明确允许使用 1 秒 timeout 作为例外。其余 timeout 只用于防止 CI 永久挂起,不作为性能阈值。 - -本计划与已有工作树对账时生产实现已经存在,无法在当前状态重放 RED-first;这里记录为实施后的回归/characterization test,不伪造 RED 结果。 - -验证: - -```bash -cargo test -p paimon vindex::range_reader::tests::response_copy_and_buffers_are_bounded_across_clones -``` - -### Task 2:把两个限制绑定为一个共享内部值(已完成) - -文件:`crates/paimon/src/vindex/range_reader.rs`。 - -1. 用一个最小的 `pub(crate) RangeReadLimiter` 绑定 shared I/O semaphore、shared response semaphore、`io_limit=C` 和 `response_limit=min(C.saturating_mul(2), Semaphore::MAX_PERMITS)`。 -2. `VindexFileReader` 持有 limiter,替换当前分离的 `permits + max_range_read_concurrency`,避免构造时传入互相矛盾的 semaphore size 和并发值。 -3. 测试构造器仍用默认 `C=32`;`try_clone_reader` clone 同一个 limiter 内的两个 `Arc`。 -4. `response_limit_saturates_at_semaphore_max_permits` 覆盖 `C > MAX_PERMITS / 2` 时不会乘法溢出或让 `Semaphore::new` panic。 - -验证:现有 `clones_share_range_read_permits` 更新为同时断言 I/O 与 response semaphore 被共享;不增加通用 limiter trait、factory 或新配置。 - -### Task 3:在两个生产入口全局共享 limiter(已完成) - -文件:`crates/paimon/src/table/vector_search_builder.rs`。 - -在以下两个现有 semaphore 创建点各创建一次 `RangeReadLimiter::new(range_read_concurrency)`,再 clone 给所有 Vindex reader: - -- `plan_and_search_pk_candidates_batch` 的 loader 路径; -- `evaluate_batch_vector_search` 的 vector entry 路径。 - -Lumina 路径保持不变。`global-index.thread-num`、native batch parallelism 和 `global-index.range-read-thread-num` 的现有语义保持不变。 - -验证:`rg -n "RangeReadLimiter::new" crates/paimon/src/table/vector_search_builder.rs` 只应命中上述两个 search scope,不能出现在 per-file closure 内。 - -### Task 4:重排 fetch pipeline(已完成) - -文件:`crates/paimon/src/vindex/range_reader.rs` 的 `fetch_range_batch`。 - -每个 range work item 按以下顺序执行: - -```text -acquire response permit - -> acquire I/O permit - -> FileRead::read - -> drop InFlightRead and I/O permit - -> validate response - -> bounded channel send - -> synchronous consume/copy - -> drop response permit -``` - -具体要求: - -1. channel capacity 从 `1` 改为 `C`。 -2. 把 `sender.send` 移进每个 range future;future 完成发送后才算 work item 完成。 -3. `buffer_unordered` 的 work window 使用 `response_limit`,真实读取仍由 I/O semaphore 限制为 `C`。 -4. 第一个 read/send 错误送达 consumer 后停止 producer;drop 剩余 stream,取消未完成 future 并释放两个 permit。 -5. `FileRead::read(...).await` 的结果先保存,再立即 drop `InFlightRead` 和 I/O permit,然后执行错误映射、长度校验和发送。 -6. 不改变 completion-order copy、range merge、scalar cache 或 stats 字段。 - -### Task 5:补齐回归测试(已完成) - -更新并保留以下已有测试语义: - -- `clones_share_range_read_permits`:两个限制均跨 clone 共享。 -- `configured_range_read_concurrency_can_exceed_32`:真实 peak I/O 仍能达到配置值。 -- `range_reads_refill_before_slowest_batch_member_finishes`:释放 I/O 后及时补位。 -- `completed_range_buffers_are_released_before_slowest_read`:消费完成的 payload 及时 drop。 -- `cloned_reader_is_not_queued_behind_an_entire_batch`:多个 reader 仍共享且公平竞争 I/O。 -- `failed_range_read_releases_permit`:改为同时验证 I/O 与 response permit,后续读取不死锁。 -- `range_io_stats_count_coalesced_reads`:`peak_in_flight_reads` 仍统计真实 read,而不是 response 生命周期。 - -运行: - -```bash -cargo fmt --all --check -cargo test -p paimon vindex::range_reader -cargo clippy -p paimon --all-targets -- -D warnings -git diff --check -``` - -本机 Apple ARM 若仍被 `paimon-vindex-core 0.3.0` 的 unstable NEON 编译问题阻塞,至少本地完成 `fmt` 和 `diff --check`,完整 test/clippy 以 Linux GitHub CI 为门禁;不得把依赖绕过改动混入本 PR。 - -2026-08-18 本地验证记录(Rust `1.97.0`): - -- `cargo fmt --all -- --check` 通过; -- `cargo test -p paimon vindex::range_reader::tests --lib`:20 passed; -- `cargo clippy -p paimon --all-targets -- -D warnings` 通过; -- `git diff --check` 通过。 - -### Task 6:复跑现有 A/B(待执行) - -用同一 benchmark 配置分别比较基线 `1e172b3` 与候选提交,每组至少 3 次取中位数: - -1. 无 local cache、无 warmup:确认冷读 QPS/平均延迟变化在约 3% 测试噪声内,Recall/NDCG 完全一致。 -2. 4 GiB memory cache、热读:候选版本至少不能比 `1e172b3` 更慢;目标是收回当前相对 `0946fcd` 约 3.4% 的 QPS 差距。 -3. 记录 peak RSS、`io_wait`、`range_permit_wait` 和 `peak_in_flight_reads`。RSS 可以高于严格 `C` 响应上限的 `1e172b3`,但不能回到保留全部 range 的增长模式;全局 `2C` 上限由确定性单测负责证明。 - -如果热 cache 仍稳定回退超过 3%,先保留双 limiter 的正确性改动但不宣称性能改善,使用现有 timing 数据确认瓶颈后再决定是否调整窗口;本计划不预先增加自适应窗口或 byte-based limiter。 - -## 5. 完成标准 - -- 真实底层 read 的 peak 不超过 `C`,响应 payload 数全局不超过 `min(2C, MAX_PERMITS)`。 -- 同步复制不再持有 I/O permit,clone/file 之间共享两种预算。 -- 错误、取消和 receiver 关闭不会泄漏 permit 或挂住后续读取。 -- Recall/NDCG 不变,冷读无稳定回退,热 cache 回退被消除或有 timing 证据解释。 -- PR 只新增上述两文件的实现/测试改动,不引入新配置、依赖或相邻重构。 From 09d550c81fb46f314cda228a33cb0217acc54a01 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 13:38:48 +0800 Subject: [PATCH 09/14] test(vindex): add range read benchmark --- crates/paimon/src/vindex/range_reader.rs | 117 +++++++++++++++++++++++ 1 file changed, 117 insertions(+) diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index edfdd1e9..59a29c47 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -610,6 +610,16 @@ mod tests { release: tokio::sync::Semaphore, } + struct BenchmarkRead { + data: Bytes, + calls: AtomicUsize, + active: AtomicUsize, + max_active: AtomicUsize, + fast_delay: Duration, + slow_every: usize, + slow_delay: Duration, + } + #[async_trait] impl FileRead for RuntimeTrackingRead { async fn read(&self, range: Range) -> crate::Result { @@ -690,6 +700,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 = DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM; + 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), + 32, + 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::>()); From ca564e67f248059a37d64d15f286dbb9c5f97daf Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 13:59:06 +0800 Subject: [PATCH 10/14] Rename Vindex range-read thread option --- crates/paimon/src/spec/core_options.rs | 4 ++-- crates/paimon/src/table/vector_search_builder.rs | 2 +- docs/src/sql.md | 4 ++-- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 05b94919..9b36a02e 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -28,7 +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_RANGE_READ_THREAD_NUM_OPTION: &str = "global-index.range-read-thread-num"; +const GLOBAL_INDEX_RANGE_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"; @@ -734,7 +734,7 @@ impl<'a> CoreOptions<'a> { } /// Maximum number of concurrent range reads shared by Vindex readers in one - /// search operation (key `global-index.range-read-thread-num`, default 32). + /// search operation (key `global-index.vindex.read-thread-num`, default 32). /// This is independent of [`Self::global_index_thread_num`]. pub fn global_index_range_read_thread_num(&self) -> crate::Result { let value = self diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index b1d023b1..28681c8c 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -3486,7 +3486,7 @@ mod tests { ); let options = HashMap::from([( - "global-index.range-read-thread-num".to_string(), + "global-index.vindex.read-thread-num".to_string(), "64".to_string(), )]); let core = CoreOptions::new(&options); diff --git a/docs/src/sql.md b/docs/src/sql.md index 58669e43..dff2cca2 100644 --- a/docs/src/sql.md +++ b/docs/src/sql.md @@ -2065,10 +2065,10 @@ deletion vectors enabled. | `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 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.range-read-thread-num` | `32` | 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.vindex.read-thread-num` | `32` | 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.range-read-thread-num` is independent of `global-index.thread-num`. +`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 `32`. From 2be6a6a8b260cc6a8d355f48342b40386a0ba9cb Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 14:03:37 +0800 Subject: [PATCH 11/14] Raise default Vindex read concurrency to 64 --- crates/paimon/src/spec/core_options.rs | 8 ++++---- crates/paimon/src/vindex/executor.rs | 2 +- docs/src/sql.md | 4 ++-- 3 files changed, 7 insertions(+), 7 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 9b36a02e..24a2c27c 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -134,7 +134,7 @@ 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_RANGE_READ_THREAD_NUM: usize = 32; +pub(crate) const DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM: usize = 64; const MAX_GLOBAL_INDEX_RANGE_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; @@ -734,7 +734,7 @@ impl<'a> CoreOptions<'a> { } /// Maximum number of concurrent range reads shared by Vindex readers in one - /// search operation (key `global-index.vindex.read-thread-num`, default 32). + /// search operation (key `global-index.vindex.read-thread-num`, default 64). /// This is independent of [`Self::global_index_thread_num`]. pub fn global_index_range_read_thread_num(&self) -> crate::Result { let value = self @@ -1557,7 +1557,7 @@ mod tests { assert_eq!(core_options.global_index_thread_num().unwrap(), 32); assert_eq!( core_options.global_index_range_read_thread_num().unwrap(), - 32 + 64 ); assert_eq!( core_options.sorted_index_records_per_range().unwrap(), @@ -1809,7 +1809,7 @@ mod tests { CoreOptions::new(&HashMap::new()) .global_index_range_read_thread_num() .unwrap(), - 32 + 64 ); for value in [32, 64] { diff --git a/crates/paimon/src/vindex/executor.rs b/crates/paimon/src/vindex/executor.rs index 173b5d7d..b48027c8 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 the default 64 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/docs/src/sql.md b/docs/src/sql.md index dff2cca2..1440cb67 100644 --- a/docs/src/sql.md +++ b/docs/src/sql.md @@ -2065,13 +2065,13 @@ deletion vectors enabled. | `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 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` | `32` | 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.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 `32`. +otherwise they use the new default of `64`. ### Variant Shredding Options From 9425619a3358694a385972f339c200982045c526 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 14:54:41 +0800 Subject: [PATCH 12/14] test(vindex): benchmark default range concurrency --- crates/paimon/src/vindex/executor.rs | 2 +- crates/paimon/src/vindex/range_reader.rs | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/crates/paimon/src/vindex/executor.rs b/crates/paimon/src/vindex/executor.rs index b48027c8..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 64 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 59a29c47..9e2d1632 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -728,7 +728,7 @@ mod tests { slow_delay: Duration, ) { const RANGE_SIZE: usize = 4 * 1024; - let concurrency = DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM; + 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]), @@ -801,7 +801,7 @@ mod tests { 256, 10, Duration::from_millis(1), - 32, + 64, Duration::from_millis(10), ) .await; From 6c37446edd6d0f158fac993e945aa6360ee418a4 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 14:55:17 +0800 Subject: [PATCH 13/14] fix(vindex): retain permits through range copies --- crates/paimon/src/vindex/range_reader.rs | 181 ++++++++--------------- 1 file changed, 63 insertions(+), 118 deletions(-) diff --git a/crates/paimon/src/vindex/range_reader.rs b/crates/paimon/src/vindex/range_reader.rs index 9e2d1632..92ad814f 100644 --- a/crates/paimon/src/vindex/range_reader.rs +++ b/crates/paimon/src/vindex/range_reader.rs @@ -143,6 +143,11 @@ impl RangeReadLimiter { } } +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. @@ -154,8 +159,6 @@ pub(crate) struct VindexFileReader { path: String, scalar_cache: Option, stats: Option>, - #[cfg(test)] - response_permit_wait_started: Option>, } impl VindexFileReader { @@ -190,8 +193,6 @@ impl VindexFileReader { path, scalar_cache: None, stats: vector_search_timing_enabled().then(|| Arc::new(RangeIoStats::default())), - #[cfg(test)] - response_permit_wait_started: None, } } @@ -240,24 +241,24 @@ 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 { + fn fetch_exact(&self, range: Range) -> io::Result { let mut result = None; - self.fetch_range_batch(std::slice::from_ref(&range), |_, data| { - result = Some(data); + self.fetch_range_batch(std::slice::from_ref(&range), |_, response| { + result = Some(response); Ok(()) })?; Ok(result.expect("one requested range")) @@ -266,7 +267,7 @@ impl VindexFileReader { fn fetch_range_batch( &self, ranges: &[Range], - mut consume: impl FnMut(usize, Bytes) -> io::Result<()>, + mut consume: impl FnMut(usize, RangeResponse) -> io::Result<()>, ) -> io::Result<()> { let reader = Arc::clone(&self.reader); let io_permits = Arc::clone(&self.limiter.io_permits); @@ -276,8 +277,6 @@ impl VindexFileReader { let stats = self.stats.clone(); let io_limit = self.limiter.io_limit; let response_limit = self.limiter.response_limit; - #[cfg(test)] - let response_permit_wait_started = self.response_permit_wait_started.clone(); let (sender, mut receiver) = tokio::sync::mpsc::channel(io_limit); self.runtime.spawn(async move { let fetched = stream::iter(requested.into_iter().enumerate().map(|(index, range)| { @@ -287,13 +286,7 @@ impl VindexFileReader { let sender = sender.clone(); let path = path.clone(); let stats = stats.clone(); - #[cfg(test)] - let response_permit_wait_started = response_permit_wait_started.clone(); async move { - #[cfg(test)] - if let Some(wait_started) = response_permit_wait_started { - wait_started.add_permits(1); - } let response_permit = response_permits .acquire_owned() .await @@ -367,13 +360,19 @@ impl VindexFileReader { 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(|| { + let (index, data, response_permit) = fetched.ok_or_else(|| { io::Error::other(format!( "vindex range read task for '{}' was cancelled", self.path )) })??; - consume(index, data) + consume( + index, + RangeResponse { + data, + _permit: response_permit, + }, + ) }); if let Some(stats) = &self.stats { stats @@ -450,14 +449,14 @@ impl VindexFileReader { .fetch_max(chunk_size, Ordering::Relaxed); } let ranges: Vec<_> = merged.iter().map(|merged| merged.range.clone()).collect(); - self.fetch_range_batch(&ranges, |merged_index, data| { + 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(&data[start..start + request.buf.len()]); + .copy_from_slice(&response.data[start..start + request.buf.len()]); } Ok(()) }) @@ -504,8 +503,6 @@ impl SeekRead for VindexFileReader { path: self.path.clone(), scalar_cache: None, stats: self.stats.clone(), - #[cfg(test)] - response_permit_wait_started: self.response_permit_wait_started.clone(), })) } @@ -603,13 +600,6 @@ mod tests { dropped: Arc, } - struct OrderedRead { - data: Bytes, - ranges: Mutex>>, - started: tokio::sync::Semaphore, - release: tokio::sync::Semaphore, - } - struct BenchmarkRead { data: Bytes, calls: AtomicUsize, @@ -673,21 +663,6 @@ mod tests { } } - #[async_trait] - impl FileRead for OrderedRead { - async fn read(&self, range: Range) -> crate::Result { - self.ranges.lock().unwrap().push(range.clone()); - self.started.add_permits(1); - acquire_test_permits( - &self.release, - 1, - "test did not release an ordered range read", - ) - .await; - Ok(self.data.slice(range.start as usize..range.end as usize)) - } - } - #[async_trait] impl FileRead for TrackingRead { async fn read(&self, range: Range) -> crate::Result { @@ -1383,6 +1358,45 @@ mod tests { ); } + #[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; @@ -1430,75 +1444,6 @@ mod tests { assert_eq!(output, vec![[1], [2], [3]]); } - #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn cloned_reader_is_not_queued_behind_an_entire_batch() { - let stride = RANGE_COALESCE_GAP + 2; - let data = Bytes::from(vec![8u8; 4 * stride as usize + 128]); - let tracking = Arc::new(OrderedRead { - data: data.clone(), - ranges: Mutex::new(Vec::new()), - started: tokio::sync::Semaphore::new(0), - release: tokio::sync::Semaphore::new(0), - }); - let source: Arc = tracking.clone(); - let mut first_reader = VindexFileReader::new_with_limiter( - source, - tokio::runtime::Handle::current(), - RangeReadLimiter::new(1), - data.len() as u64, - "index".to_string(), - ); - let mut second_reader = first_reader.try_clone_reader().unwrap().unwrap(); - let response_wait_started = Arc::new(tokio::sync::Semaphore::new(0)); - second_reader.response_permit_wait_started = Some(Arc::clone(&response_wait_started)); - - let first = tokio::task::spawn_blocking(move || { - let mut first = [0u8; 1]; - let mut second = [0u8; 1]; - let mut third = [0u8; 1]; - first_reader - .pread(&mut [ - ReadRequest::new(0, &mut first), - ReadRequest::new(stride, &mut second), - ReadRequest::new(2 * stride, &mut third), - ]) - .unwrap(); - }); - acquire_test_permits(&tracking.started, 1, "first reader did not start").await; - - let second = tokio::task::spawn_blocking(move || { - let mut output = [0u8; 128]; - second_reader - .pread(&mut [ReadRequest::new(3 * stride, &mut output)]) - .unwrap(); - }); - acquire_test_permits( - &response_wait_started, - 1, - "cloned reader did not start waiting for a response permit", - ) - .await; - - tracking.release.add_permits(1); - acquire_test_permits(&tracking.started, 1, "second range read did not start").await; - tracking.release.add_permits(1); - acquire_test_permits(&tracking.started, 1, "cloned range read did not start").await; - let first_three_ranges = tracking.ranges.lock().unwrap()[..3].to_vec(); - tracking.release.add_permits(3); - tokio::time::timeout(Duration::from_secs(5), first) - .await - .expect("first reader task did not finish") - .unwrap(); - tokio::time::timeout(Duration::from_secs(5), second) - .await - .expect("cloned reader task did not finish") - .unwrap(); - - assert!(first_three_ranges - .iter() - .any(|range| range.start == 3 * stride)); - } - #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn failed_range_read_releases_permit() { let data = Bytes::from(vec![9u8; 1024]); From be51376fa3acb8a0622401eedd98028eaba0d610 Mon Sep 17 00:00:00 2001 From: yantian Date: Tue, 18 Aug 2026 15:28:46 +0800 Subject: [PATCH 14/14] test(vindex): update default concurrency assertions --- crates/paimon/src/table/vector_search_builder.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 28681c8c..23c7eefb 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -3478,20 +3478,20 @@ mod tests { let default_core = CoreOptions::new(&default_options); assert_eq!( vindex_concurrency_limits(&default_core, 1, 32).unwrap(), - (1, 32) + (1, 64) ); assert_eq!( vindex_concurrency_limits(&default_core, 8, 4).unwrap(), - (4, 32) + (4, 64) ); let options = HashMap::from([( "global-index.vindex.read-thread-num".to_string(), - "64".to_string(), + "48".to_string(), )]); let core = CoreOptions::new(&options); - assert_eq!(vindex_concurrency_limits(&core, 1, 32).unwrap(), (1, 64)); - assert_eq!(vindex_concurrency_limits(&core, 8, 4).unwrap(), (4, 64)); + 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 {