This is an automated email from the ASF dual-hosted git repository.
XiaoHongbo-Hope 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 2ef07025 perf(table): search global index shards concurrently (#589)
2ef07025 is described below
commit 2ef070251759403ae9e09451922eae287c0852a8
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 23 00:14:02 2026 +0800
perf(table): search global index shards concurrently (#589)
---
crates/paimon/src/spec/core_options.rs | 8 +-
.../src/table/btree_global_index_build_builder.rs | 2 +
crates/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 43230ebf..eec1ca57 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<usize> {
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 07bbe166..b97204a7 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(),
@@ -1359,6 +1360,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 b6f50e88..c0122d44 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<dyn Fn(&[u8], &[u8]) -> 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<T, Fut, Acc, Fold>(
+ futures: impl IntoIterator<Item = Fut>,
+ max_concurrency: usize,
+ mut accumulator: Acc,
+ mut fold: Fold,
+) -> Result<Acc>
+where
+ Fut: Future<Output = Result<T>>,
+ 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<RowRange>,
evaluated_field_ids: HashSet<i32>,
@@ -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<BoxedCmp>),
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<usize>,
+}
+
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<Option<Self>> {
+ 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<i32,
Vec<GlobalIndexEntry>> =
std::collections::HashMap::new();
let mut coverage_by_field: HashMap<i32, Vec<RowRange>> =
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<Option<RoaringTreemap>> {
+ 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<i64>,
@@ -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::<Vec<_>>());
+ 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(),