Skip to content
Open
87 changes: 83 additions & 4 deletions crates/paimon/src/spec/core_options.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ const VECTOR_INDEX_SEARCH_MODE_OPTION: &str = "vector-index.search-mode";
const FULL_TEXT_INDEX_SEARCH_MODE_OPTION: &str = "full-text-index.search-mode";
const GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION: &str = "global-index.row-count-per-shard";
const GLOBAL_INDEX_THREAD_NUM_OPTION: &str = "global-index.thread-num";
const GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION: &str = "global-index.range-read-thread-num";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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";
Expand Down Expand Up @@ -133,6 +134,8 @@ const DYNAMIC_BUCKET_TARGET_ROW_NUM_OPTION: &str = "dynamic-bucket.target-row-nu
const DEFAULT_DYNAMIC_BUCKET_TARGET_ROW_NUM: i64 = 200_000;
const DEFAULT_GLOBAL_INDEX_ROW_COUNT_PER_SHARD: i64 = 100_000;
const DEFAULT_GLOBAL_INDEX_THREAD_NUM: i64 = 32;
pub(crate) const DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM: usize = 32;
const MAX_GLOBAL_INDEX_RANGE_READ_THREAD_NUM: i64 = tokio::sync::Semaphore::MAX_PERMITS as i64;
const MAX_GLOBAL_INDEX_THREAD_NUM: i64 = {
let tokio_max = (usize::MAX >> 3) as u64;
let i32_max = i32::MAX as u64;
Expand Down Expand Up @@ -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_range_read_thread_num`] instead. A value of `1`
/// makes these search tasks sequential, but does not serialize Vindex range
/// reads. A non-positive value, or one above [`MAX_GLOBAL_INDEX_THREAD_NUM`],
/// is a misconfiguration and fails loud rather than being silently clamped.
pub fn global_index_thread_num(&self) -> crate::Result<usize> {
let value = self
.parse_i64_option(GLOBAL_INDEX_THREAD_NUM_OPTION)?
Expand All @@ -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.range-read-thread-num`, default 32).
/// This is independent of [`Self::global_index_thread_num`].
pub fn global_index_range_read_thread_num(&self) -> crate::Result<usize> {
let value = self
.parse_i64_option(GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION)?
.unwrap_or(DEFAULT_GLOBAL_INDEX_RANGE_READ_THREAD_NUM as i64);
if value <= 0 {
return Err(crate::Error::DataInvalid {
message: format!(
"Option '{}' must be greater than 0, got: {}",
GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION, value
),
source: None,
});
}
if value > MAX_GLOBAL_INDEX_RANGE_READ_THREAD_NUM {
return Err(crate::Error::DataInvalid {
message: format!(
"Option '{}' must not exceed {}, got: {}",
GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION,
MAX_GLOBAL_INDEX_RANGE_READ_THREAD_NUM,
value
),
source: None,
});
}
Ok(value as usize)
}

pub fn sorted_index_records_per_range(&self) -> crate::Result<i64> {
let value = self
.parse_i64_option(SORTED_INDEX_RECORDS_PER_RANGE_OPTION)?
Expand Down Expand Up @@ -1520,6 +1555,10 @@ mod tests {
100_000
);
assert_eq!(core_options.global_index_thread_num().unwrap(), 32);
assert_eq!(
core_options.global_index_range_read_thread_num().unwrap(),
32
);
assert_eq!(
core_options.sorted_index_records_per_range().unwrap(),
100_000
Expand Down Expand Up @@ -1764,6 +1803,46 @@ mod tests {
);
}

#[test]
fn test_global_index_range_read_thread_num_default_and_custom() {
assert_eq!(
CoreOptions::new(&HashMap::new())
.global_index_range_read_thread_num()
.unwrap(),
32
);

for value in [32, 64] {
let options = HashMap::from([(
GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION.to_string(),
value.to_string(),
)]);
assert_eq!(
CoreOptions::new(&options)
.global_index_range_read_thread_num()
.unwrap(),
value
);
}
}

#[test]
fn test_global_index_range_read_thread_num_rejects_invalid_values() {
for value in [
"0".to_string(),
"abc".to_string(),
(MAX_GLOBAL_INDEX_RANGE_READ_THREAD_NUM + 1).to_string(),
] {
let options =
HashMap::from([(GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION.to_string(), value)]);
let err = CoreOptions::new(&options)
.global_index_range_read_thread_num()
.expect_err("invalid range-read-thread-num should fail");
assert!(matches!(err, crate::Error::DataInvalid { message, .. }
if message.contains(GLOBAL_INDEX_RANGE_READ_THREAD_NUM_OPTION)));
}
}

#[test]
fn test_sorted_index_records_per_range_rejects_invalid_values() {
for value in ["0", "-1", "abc"] {
Expand Down
93 changes: 67 additions & 26 deletions crates/paimon/src/table/vector_search_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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,
Expand All @@ -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_range_read_thread_num()?,
))
}

pub struct VectorSearchBuilder<'a> {
table: &'a Table,
vector_column: Option<String>,
Expand Down Expand Up @@ -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`
Expand All @@ -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 {
Expand All @@ -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,
),
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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(),
);
Expand Down Expand Up @@ -3446,11 +3473,25 @@ mod tests {
}

#[test]
fn vindex_batch_parallelism_tracks_active_entries() {
assert_eq!(vindex_index_parallelism(1, 1), 1);
assert_eq!(vindex_index_parallelism(1, 64), 1);
assert_eq!(vindex_index_parallelism(8, 4), 4);
assert_eq!(vindex_index_parallelism(4, 8), 4);
fn vindex_concurrency_limits_are_independent() {
let default_options = HashMap::new();
let default_core = CoreOptions::new(&default_options);
assert_eq!(
vindex_concurrency_limits(&default_core, 1, 32).unwrap(),
(1, 32)
);
assert_eq!(
vindex_concurrency_limits(&default_core, 8, 4).unwrap(),
(4, 32)
);

let options = HashMap::from([(
"global-index.range-read-thread-num".to_string(),
"64".to_string(),
)]);
let core = CoreOptions::new(&options);
assert_eq!(vindex_concurrency_limits(&core, 1, 32).unwrap(), (1, 64));
assert_eq!(vindex_concurrency_limits(&core, 8, 4).unwrap(), (4, 64));
}

fn make_field(id: i32, name: &str) -> DataField {
Expand Down
Loading
Loading