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