This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new e4959719 perf(vindex): decouple vector read threads and remove chunk 
barrier (#720)
e4959719 is described below

commit e4959719fecf509056ffe39dd6e39f6e0396e12c
Author: jerry <[email protected]>
AuthorDate: Tue Aug 18 18:00:19 2026 +0800

    perf(vindex): decouple vector read threads and remove chunk barrier (#720)
---
 crates/paimon/src/spec/core_options.rs           |  89 ++-
 crates/paimon/src/table/vector_search_builder.rs |  93 ++-
 crates/paimon/src/vindex/executor.rs             |   2 +-
 crates/paimon/src/vindex/range_reader.rs         | 835 ++++++++++++++++++++---
 docs/src/sql.md                                  |   8 +-
 5 files changed, 902 insertions(+), 125 deletions(-)

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

Reply via email to