From 79a44a472356dd2b7175fd3663c64b80545208d6 Mon Sep 17 00:00:00 2001 From: XiaoHongbo Date: Wed, 22 Jul 2026 07:40:30 -0700 Subject: [PATCH] perf(table): search global index shards concurrently --- crates/paimon/src/spec/core_options.rs | 8 +- .../table/btree_global_index_build_builder.rs | 2 + .../paimon/src/table/global_index_scanner.rs | 317 ++++++++++++++---- crates/paimon/src/table/table_scan.rs | 14 +- 4 files changed, 260 insertions(+), 81 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index dd545f19..4c8ae18d 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -578,10 +578,10 @@ impl<'a> CoreOptions<'a> { /// Maximum number of concurrent tasks for global-index I/O, mirroring Java /// `CoreOptions.GLOBAL_INDEX_THREAD_NUM` (key `global-index.thread-num`, - /// default 32). Used as the fan-out limit for the primary-key vector search - /// (per-bucket and per-exact-file). A value of `1` reproduces strict - /// sequential execution. A non-positive value is a misconfiguration and fails - /// loud rather than being silently clamped. + /// default 32). Used as the per-operation fan-out limit for sorted BTree and + /// bitmap shard reads and for primary-key vector search. A value of `1` + /// reproduces strict sequential execution. A non-positive value is a + /// misconfiguration and fails loud rather than being silently clamped. pub fn global_index_thread_num(&self) -> crate::Result { let value = self .parse_i64_option(GLOBAL_INDEX_THREAD_NUM_OPTION)? diff --git a/crates/paimon/src/table/btree_global_index_build_builder.rs b/crates/paimon/src/table/btree_global_index_build_builder.rs index 0b759b1d..b133f3a1 100644 --- a/crates/paimon/src/table/btree_global_index_build_builder.rs +++ b/crates/paimon/src/table/btree_global_index_build_builder.rs @@ -1231,6 +1231,7 @@ mod tests { predicates: &[predicate], schema_fields: table.schema().fields(), search_mode: GlobalIndexSearchMode::Fast, + global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, bitmap_fallback_scan_max_size: i64::MAX, next_row_id: snapshot.next_row_id(), @@ -1332,6 +1333,7 @@ mod tests { predicates: &[predicate], schema_fields: table.schema().fields(), search_mode: GlobalIndexSearchMode::Fast, + global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, bitmap_fallback_scan_max_size: i64::MAX, next_row_id: snapshot.next_row_id(), diff --git a/crates/paimon/src/table/global_index_scanner.rs b/crates/paimon/src/table/global_index_scanner.rs index c535db7f..b930c081 100644 --- a/crates/paimon/src/table/global_index_scanner.rs +++ b/crates/paimon/src/table/global_index_scanner.rs @@ -24,7 +24,7 @@ use super::bitmap_global_index_reader::BitmapGlobalIndexReader; use super::global_index_types::{ normalize_sorted_global_index_type, BITMAP_GLOBAL_INDEX_TYPE, BTREE_GLOBAL_INDEX_TYPE, }; -use crate::btree::query::{extract_between, IndexQuery}; +use crate::btree::query::{extract_between, BetweenInfo, IndexQuery}; use crate::btree::{make_key_comparator, serialize_datum, BTreeIndexMeta, BTreeIndexReader}; use crate::deletion_vector::DeletionVectorFactory; use crate::io::FileIO; @@ -34,9 +34,11 @@ use crate::spec::{ }; use crate::table::{DeletionFile, RowRange, Table}; use crate::{Error, Result}; +use futures::{StreamExt, TryStreamExt}; use roaring::RoaringTreemap; use std::cmp::Ordering; use std::collections::{HashMap, HashSet}; +use std::future::Future; use std::sync::Mutex; type BoxedCmp = Box Ordering + Send + Sync>; @@ -50,6 +52,25 @@ type PredicateTuple<'a> = (PredicateOperator, &'a [Datum], &'a DataType); const DELETION_VECTORS_INDEX_TYPE: &str = "DELETION_VECTORS"; const INDEX_DIR: &str = "index"; +async fn try_fold_bounded( + futures: impl IntoIterator, + max_concurrency: usize, + mut accumulator: Acc, + mut fold: Fold, +) -> Result +where + Fut: Future>, + Fold: FnMut(&mut Acc, T), +{ + debug_assert!(max_concurrency > 0); + let stream = futures::stream::iter(futures).buffer_unordered(max_concurrency); + futures::pin_mut!(stream); + while let Some(value) = stream.try_next().await? { + fold(&mut accumulator, value); + } + Ok(accumulator) +} + struct GlobalIndexScanResult { row_ranges: Vec, evaluated_field_ids: HashSet, @@ -64,6 +85,7 @@ struct GlobalIndexScanResult { pub(crate) struct GlobalIndexScanner { file_io: FileIO, table_path: String, + global_index_thread_num: usize, btree_fallback_scan_max_size: i64, bitmap_fallback_scan_max_size: i64, /// Global index entries grouped by field_id. @@ -91,6 +113,15 @@ enum GlobalIndexFileKind { Bitmap, } +impl GlobalIndexFileKind { + fn name(self) -> &'static str { + match self { + Self::BTree => "BTree", + Self::Bitmap => "bitmap", + } + } +} + enum OpenedGlobalIndexReader { BTree(BTreeIndexReader), Bitmap(BitmapGlobalIndexReader), @@ -104,6 +135,13 @@ struct FallbackScanPlan { allow_bitmap: bool, } +struct EntryQueryPlan { + entry_idx: usize, + between_matches: bool, + between_evaluated: bool, + matching_predicates: Vec, +} + impl FallbackScanPlan { fn allowed(self, kind: GlobalIndexFileKind) -> bool { match kind { @@ -155,11 +193,18 @@ impl GlobalIndexScanner { pub(crate) fn create( file_io: &FileIO, table_path: &str, + global_index_thread_num: usize, btree_fallback_scan_max_size: i64, bitmap_fallback_scan_max_size: i64, index_entries: &[IndexManifestEntry], schema_fields: &[DataField], ) -> Result> { + if global_index_thread_num == 0 { + return Err(Error::DataInvalid { + message: "Global index thread count must be greater than 0".to_string(), + source: None, + }); + } let mut entries_by_field: std::collections::HashMap> = std::collections::HashMap::new(); let mut coverage_by_field: HashMap> = HashMap::new(); @@ -243,6 +288,7 @@ impl GlobalIndexScanner { Ok(Some(Self { file_io: file_io.clone(), table_path: table_path.trim_end_matches('/').to_string(), + global_index_thread_num, btree_fallback_scan_max_size, bitmap_fallback_scan_max_size, entries_by_field: entries_by_field.into_iter().collect(), @@ -394,8 +440,6 @@ impl GlobalIndexScanner { predicates }; - let mut all_row_ids = RoaringTreemap::new(); - // Pre-compute comparators and serialized keys for file-level pruning per predicate let pruning_info: Vec<_> = effective_predicates .iter() @@ -443,6 +487,7 @@ impl GlobalIndexScanner { .as_ref() .map(|_| self.fallback_scan_plan(entries, &between_matches_by_entry)); + let mut query_plans = Vec::with_capacity(entries.len()); for (entry_idx, entry) in entries.iter().enumerate() { // Also check if between range may match let between_matches = between @@ -500,83 +545,120 @@ impl GlobalIndexScanner { continue; } - let data_type = between - .as_ref() - .map(|b| b.data_type) - .or_else(|| effective_predicates.first().map(|p| p.2)) - .unwrap_or(predicates[0].2); - let mut reader = if (between_matches && between_evaluated_for_entry) - || !matching_predicates.is_empty() - { - Some( - self.open_reader_for_entry(entry, &entry.meta, data_type) - .await?, - ) - } else { - None - }; + query_plans.push(EntryQueryPlan { + entry_idx, + between_matches, + between_evaluated: between_evaluated_for_entry, + matching_predicates, + }); + } - let mut file_result = None; - - // Execute between query first if applicable - if between_matches && between_evaluated_for_entry { - if let Some(b) = &between { - let from_key = serialize_datum(b.from, b.data_type); - let to_key = serialize_datum(b.to, b.data_type); - let bitmap = reader - .as_ref() - .expect("reader is opened when between matches") - .range_query( - &from_key, - &to_key, - b.data_type, - b.from_inclusive, - b.to_inclusive, - ) - .await - .map_err(|e| crate::Error::DataInvalid { - message: "Global index query failed".to_string(), - source: Some(Box::new(e)), - })?; - file_result = Some(bitmap); + // Complete all pruning and fallback decisions before starting shard I/O. + // A later unsupported shard must fall back to the normal scan without an + // earlier shard racing it with an I/O or query error. + let data_type = between + .as_ref() + .map(|b| b.data_type) + .or_else(|| effective_predicates.first().map(|p| p.2)) + .unwrap_or(predicates[0].2); + let between = between.as_ref(); + let futures = query_plans.into_iter().map(|plan| async move { + let entry = &entries[plan.entry_idx]; + let result = self + .query_entry(entry, data_type, between, &plan, effective_predicates) + .await?; + Ok((entry.row_range_start, result)) + }); + let all_row_ids = try_fold_bounded( + futures, + self.global_index_thread_num, + RoaringTreemap::new(), + |all_row_ids, (row_range_start, file_result)| { + if let Some(bitmap) = file_result { + for row_id in bitmap.iter() { + all_row_ids.insert(row_id + row_range_start as u64); + } } - } + }, + ) + .await?; - // Evaluate remaining predicates - for &idx in &matching_predicates { - let (op, literals, dt) = &effective_predicates[idx]; - let bitmap = reader - .as_ref() - .expect("reader is opened when predicates match") - .query(*op, literals, dt) - .await - .map_err(|e| crate::Error::DataInvalid { - message: "Global index query failed".to_string(), - source: Some(Box::new(e)), - })?; - file_result = Some(match file_result { - None => bitmap, - Some(mut existing) => { - existing &= bitmap; - existing - } - }); - } + Ok(Some(bitmap_to_ranges(&all_row_ids))) + } - // Return BTree readers to cache. Bitmap readers are cheap wrappers - // around one opened file and are not cached yet. - if let Some(OpenedGlobalIndexReader::BTree(reader)) = reader.take() { - self.return_reader(entry.file_name.clone(), reader); - } + async fn query_entry( + &self, + entry: &GlobalIndexEntry, + data_type: &DataType, + between: Option<&BetweenInfo<'_>>, + plan: &EntryQueryPlan, + effective_predicates: &[(PredicateOperator, &[Datum], &DataType)], + ) -> Result> { + let mut reader = if (plan.between_matches && plan.between_evaluated) + || !plan.matching_predicates.is_empty() + { + Some( + self.open_reader_for_entry(entry, &entry.meta, data_type) + .await?, + ) + } else { + None + }; + let mut file_result = None; - if let Some(bitmap) = file_result { - for rid in bitmap.iter() { - all_row_ids.insert(rid + entry.row_range_start as u64); + if plan.between_matches && plan.between_evaluated { + let between = between.expect("evaluated between query is present"); + let from_key = serialize_datum(between.from, between.data_type); + let to_key = serialize_datum(between.to, between.data_type); + let bitmap = reader + .as_ref() + .expect("reader is opened when between matches") + .range_query( + &from_key, + &to_key, + between.data_type, + between.from_inclusive, + between.to_inclusive, + ) + .await + .map_err(|error| Self::query_error(entry, error))?; + file_result = Some(bitmap); + } + + for &idx in &plan.matching_predicates { + let (op, literals, data_type) = &effective_predicates[idx]; + let bitmap = reader + .as_ref() + .expect("reader is opened when predicates match") + .query(*op, literals, data_type) + .await + .map_err(|error| Self::query_error(entry, error))?; + file_result = Some(match file_result { + None => bitmap, + Some(mut existing) => { + existing &= bitmap; + existing } - } + }); } - Ok(Some(bitmap_to_ranges(&all_row_ids))) + // Each concurrent task owns its reader. Only return it to the shared + // cache after all predicates for this shard have completed. + if let Some(OpenedGlobalIndexReader::BTree(reader)) = reader.take() { + self.return_reader(entry.file_name.clone(), reader); + } + Ok(file_result) + } + + fn query_error(entry: &GlobalIndexEntry, error: std::io::Error) -> Error { + Error::DataInvalid { + message: format!( + "Global index query failed for {} file '{}'", + entry.index_type.name(), + entry.file_name + ), + source: Some(Box::new(error)), + } } /// Get a cached reader or open a new one for the given file. @@ -1208,6 +1290,7 @@ pub(crate) struct GlobalIndexEvaluation<'a> { pub(crate) predicates: &'a [Predicate], pub(crate) schema_fields: &'a [DataField], pub(crate) search_mode: GlobalIndexSearchMode, + pub(crate) global_index_thread_num: usize, pub(crate) btree_fallback_scan_max_size: i64, pub(crate) bitmap_fallback_scan_max_size: i64, pub(crate) next_row_id: Option, @@ -1220,6 +1303,7 @@ pub(crate) async fn evaluate_global_index( let scanner = match GlobalIndexScanner::create( evaluation.file_io, evaluation.table_path, + evaluation.global_index_thread_num, evaluation.btree_fallback_scan_max_size, evaluation.bitmap_fallback_scan_max_size, evaluation.index_entries, @@ -1248,6 +1332,37 @@ pub(crate) async fn evaluate_global_index( #[cfg(test)] mod tests { use super::*; + use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}; + use std::sync::Arc; + + #[tokio::test] + async fn test_try_fold_bounded_respects_concurrency_limit() { + for limit in [1, 3] { + let active = Arc::new(AtomicUsize::new(0)); + let peak = Arc::new(AtomicUsize::new(0)); + let futures = (0..9usize).map(|value| { + let active = Arc::clone(&active); + let peak = Arc::clone(&peak); + async move { + let current = active.fetch_add(1, AtomicOrdering::SeqCst) + 1; + peak.fetch_max(current, AtomicOrdering::SeqCst); + tokio::task::yield_now().await; + active.fetch_sub(1, AtomicOrdering::SeqCst); + Ok::<_, crate::Error>(value) + } + }); + + let mut values = try_fold_bounded(futures, limit, Vec::new(), |values, value| { + values.push(value) + }) + .await + .unwrap(); + values.sort_unstable(); + + assert_eq!(values, (0..9).collect::>()); + assert_eq!(peak.load(AtomicOrdering::SeqCst), limit); + } + } #[test] fn test_bitmap_to_ranges() { @@ -1522,6 +1637,7 @@ mod tests { predicates, schema_fields: fields, search_mode: GlobalIndexSearchMode::Fast, + global_index_thread_num: 32, btree_fallback_scan_max_size, bitmap_fallback_scan_max_size, next_row_id: None, @@ -1564,6 +1680,7 @@ mod tests { let scanner = GlobalIndexScanner::create( &file_io, "memory:/t", + 32, i64::MAX, i64::MAX, &entries, @@ -1592,6 +1709,7 @@ mod tests { let scanner = GlobalIndexScanner::create( &file_io, "memory:/t", + 32, i64::MAX, i64::MAX, &entries, @@ -1620,6 +1738,7 @@ mod tests { let scanner = GlobalIndexScanner::create( &file_io, "memory:/t", + 32, i64::MAX, i64::MAX, &entries, @@ -1655,6 +1774,7 @@ mod tests { let scanner = GlobalIndexScanner::create( &file_io, "memory:/t", + 32, i64::MAX, i64::MAX, &entries, @@ -1679,6 +1799,7 @@ mod tests { let scanner = GlobalIndexScanner::create( &file_io, "memory:/t", + 32, i64::MAX, i64::MAX, &entries, @@ -1709,6 +1830,7 @@ mod tests { let scanner = GlobalIndexScanner::create( &file_io, "memory:/t", + 32, i64::MAX, i64::MAX, &[entry], @@ -2225,6 +2347,54 @@ mod tests { ); } + #[tokio::test] + async fn test_fallback_preflight_happens_before_shard_io() { + let file_io = crate::io::FileIOBuilder::new("memory").build().unwrap(); + let table_path = "memory:/missing-index-files"; + let meta = BTreeIndexMeta::new(Some(b"a".to_vec()), Some(b"z".to_vec()), false); + let fields = string_schema_fields(); + let predicates = vec![Predicate::Leaf { + column: "name".to_string(), + index: 0, + data_type: DataType::VarChar(crate::spec::VarCharType::string_type()), + op: PredicateOperator::Contains, + literals: vec![Datum::String("middle".to_string())], + }]; + + let mut btree = make_global_index_entry_with_type( + BTREE_GLOBAL_INDEX_TYPE, + "missing-btree.index", + 1, + 0, + 99, + &meta, + ); + btree.index_file.file_size = 1; + let mut bitmap = make_global_index_entry_with_type( + BITMAP_GLOBAL_INDEX_TYPE, + "missing-bitmap.index", + 1, + 100, + 199, + &meta, + ); + bitmap.index_file.file_size = 1; + + let result = evaluate_global_index_fast_with_fallback_size( + &file_io, + table_path, + &[btree, bitmap], + &predicates, + &fields, + 1, + 0, + ) + .await + .expect("fallback must be decided before opening an earlier shard"); + + assert!(result.is_none()); + } + #[tokio::test] async fn test_evaluate_global_index_full_mode_includes_unindexed_tail() { let (file_io, table_path, file_name, _tmp) = @@ -2241,6 +2411,7 @@ mod tests { predicates: &predicates, schema_fields: &fields, search_mode: GlobalIndexSearchMode::Full, + global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, bitmap_fallback_scan_max_size: i64::MAX, next_row_id: Some(150), @@ -2293,6 +2464,7 @@ mod tests { predicates: &predicates, schema_fields: &fields, search_mode: GlobalIndexSearchMode::Full, + global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, bitmap_fallback_scan_max_size: i64::MAX, next_row_id: Some(100), @@ -2326,6 +2498,7 @@ mod tests { predicates: &predicates, schema_fields: &fields, search_mode: GlobalIndexSearchMode::Detail, + global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, bitmap_fallback_scan_max_size: i64::MAX, next_row_id: Some(150), diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 8044590c..af36b177 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1325,17 +1325,20 @@ impl<'a> PaimonTableScan<'a> { entry.1.push(file); } - let global_index_search_mode = if data_evolution_enabled + let global_index_settings = if data_evolution_enabled && core_options.global_index_enabled() && !self.data_predicates.is_empty() { - Some(core_options.global_index_search_mode()?) + Some(( + core_options.global_index_search_mode()?, + core_options.global_index_thread_num()?, + )) } else { None }; let global_index_detail_data_ranges = if matches!( - global_index_search_mode, - Some(GlobalIndexSearchMode::Detail) + global_index_settings, + Some((GlobalIndexSearchMode::Detail, _)) ) { global_index_detail_data_ranges(&groups) } else { @@ -1390,7 +1393,7 @@ impl<'a> PaimonTableScan<'a> { // Use pushed-down row_ranges first; otherwise try global index. let row_ranges = if self.row_ranges.is_some() { self.row_ranges.clone() - } else if let Some(search_mode) = global_index_search_mode { + } else if let Some((search_mode, global_index_thread_num)) = global_index_settings { super::global_index_scanner::evaluate_global_index( super::global_index_scanner::GlobalIndexEvaluation { file_io, @@ -1399,6 +1402,7 @@ impl<'a> PaimonTableScan<'a> { predicates: &self.data_predicates, schema_fields: self.table.schema().fields(), search_mode, + global_index_thread_num, btree_fallback_scan_max_size: btree_index_fallback_scan_max_size, bitmap_fallback_scan_max_size: bitmap_index_fallback_scan_max_size, next_row_id: snapshot.next_row_id(),