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 792c9d9c perf(table): search global index fields concurrently (#610)
792c9d9c is described below

commit 792c9d9cf63b99aa3436d9cd64dc290b5b4229bf
Author: QuakeWang <[email protected]>
AuthorDate: Sun Jul 26 19:40:49 2026 +0800

    perf(table): search global index fields concurrently (#610)
---
 crates/paimon/src/table/global_index_scanner.rs | 180 +++++++++++++++++++++---
 1 file changed, 163 insertions(+), 17 deletions(-)

diff --git a/crates/paimon/src/table/global_index_scanner.rs 
b/crates/paimon/src/table/global_index_scanner.rs
index ad0fa022..80cf6c68 100644
--- a/crates/paimon/src/table/global_index_scanner.rs
+++ b/crates/paimon/src/table/global_index_scanner.rs
@@ -42,7 +42,11 @@ use roaring::RoaringTreemap;
 use std::cmp::Ordering;
 use std::collections::{HashMap, HashSet};
 use std::future::Future;
-use std::sync::Mutex;
+use std::sync::{Arc, Mutex};
+use tokio::sync::Semaphore;
+
+#[cfg(test)]
+use std::sync::atomic::{AtomicUsize as TestAtomicUsize, Ordering as 
TestOrdering};
 
 type BoxedCmp = Box<dyn Fn(&[u8], &[u8]) -> Ordering + Send + Sync>;
 
@@ -74,6 +78,40 @@ where
     Ok(accumulator)
 }
 
+#[cfg(test)]
+#[derive(Default)]
+struct QueryIoProbe {
+    active: TestAtomicUsize,
+    peak: TestAtomicUsize,
+}
+
+#[cfg(test)]
+impl QueryIoProbe {
+    async fn enter(&self) -> QueryIoProbeGuard<'_> {
+        let current = self.active.fetch_add(1, TestOrdering::SeqCst) + 1;
+        self.peak.fetch_max(current, TestOrdering::SeqCst);
+        let guard = QueryIoProbeGuard { probe: self };
+        tokio::task::yield_now().await;
+        guard
+    }
+
+    fn peak(&self) -> usize {
+        self.peak.load(TestOrdering::SeqCst)
+    }
+}
+
+#[cfg(test)]
+struct QueryIoProbeGuard<'a> {
+    probe: &'a QueryIoProbe,
+}
+
+#[cfg(test)]
+impl Drop for QueryIoProbeGuard<'_> {
+    fn drop(&mut self) {
+        self.probe.active.fetch_sub(1, TestOrdering::SeqCst);
+    }
+}
+
 struct GlobalIndexScanResult {
     row_ranges: Vec<RowRange>,
     evaluated_field_ids: HashSet<i32>,
@@ -89,6 +127,8 @@ pub(crate) struct GlobalIndexScanner {
     file_io: FileIO,
     table_path: String,
     global_index_thread_num: usize,
+    /// Scan-scoped shard I/O budget shared by all indexed fields.
+    query_semaphore: Arc<Semaphore>,
     btree_fallback_scan_max_size: i64,
     bitmap_fallback_scan_max_size: i64,
     /// Global index entries grouped by field_id.
@@ -99,6 +139,8 @@ pub(crate) struct GlobalIndexScanner {
     schema_fields: Vec<DataField>,
     /// Cache of opened BTree readers, keyed by file name.
     reader_cache: Mutex<HashMap<String, BTreeIndexReader<BoxedCmp>>>,
+    #[cfg(test)]
+    query_io_probe: Option<Arc<QueryIoProbe>>,
 }
 
 /// A resolved global index entry with parsed metadata.
@@ -242,6 +284,15 @@ impl GlobalIndexScanner {
                 source: None,
             });
         }
+        if global_index_thread_num > Semaphore::MAX_PERMITS {
+            return Err(Error::DataInvalid {
+                message: format!(
+                    "Global index thread count must not exceed {}",
+                    Semaphore::MAX_PERMITS
+                ),
+                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();
@@ -326,12 +377,15 @@ impl GlobalIndexScanner {
             file_io: file_io.clone(),
             table_path: table_path.trim_end_matches('/').to_string(),
             global_index_thread_num,
+            query_semaphore: Arc::new(Semaphore::new(global_index_thread_num)),
             btree_fallback_scan_max_size,
             bitmap_fallback_scan_max_size,
             entries_by_field: entries_by_field.into_iter().collect(),
             coverage_by_field,
             schema_fields: schema_fields.to_vec(),
             reader_cache: Mutex::new(HashMap::new()),
+            #[cfg(test)]
+            query_io_probe: None,
         }))
     }
 
@@ -399,21 +453,35 @@ impl GlobalIndexScanner {
                         non_leaf_children.push(child);
                     }
 
-                    let mut row_ranges: Option<Vec<RowRange>> = None;
-                    let mut evaluated_field_ids = HashSet::new();
-
-                    // Evaluate grouped leaves (one open per file)
+                    // Evaluate independent fields concurrently while keeping 
predicates for the
+                    // same field together so each index file is opened only 
once.
+                    let mut leaf_futures = 
Vec::with_capacity(leaf_groups.len());
                     for (field_id, predicates) in &leaf_groups {
                         if let Some(entries) = 
self.entries_for_field(*field_id) {
-                            if let Some(ranges) = self.evaluate_leaf(entries, 
predicates).await? {
-                                row_ranges = Some(match row_ranges {
+                            let field_id = *field_id;
+                            let predicates = predicates.as_slice();
+                            leaf_futures.push(async move {
+                                let ranges = self.evaluate_leaf(entries, 
predicates).await?;
+                                Ok((field_id, ranges))
+                            });
+                        }
+                    }
+                    let leaf_group_count = leaf_futures.len();
+                    let (mut row_ranges, mut evaluated_field_ids) = 
try_fold_bounded(
+                        leaf_futures,
+                        leaf_group_count.max(1),
+                        (None::<Vec<RowRange>>, HashSet::new()),
+                        |(row_ranges, evaluated_field_ids), (field_id, 
ranges)| {
+                            if let Some(ranges) = ranges {
+                                *row_ranges = Some(match row_ranges.take() {
                                     None => ranges,
                                     Some(existing) => 
intersect_sorted_ranges(&existing, &ranges),
                                 });
-                                evaluated_field_ids.insert(*field_id);
+                                evaluated_field_ids.insert(field_id);
                             }
-                        }
-                    }
+                        },
+                    )
+                    .await?;
 
                     // Evaluate non-leaf children recursively
                     for child in non_leaf_children {
@@ -640,13 +708,25 @@ impl GlobalIndexScanner {
             .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 futures =
+            query_plans.into_iter().map(|plan| async move {
+                let entry = &entries[plan.entry_idx];
+                let _permit = 
self.query_semaphore.acquire().await.map_err(|error| {
+                    Error::UnexpectedError {
+                        message: "global-index query concurrency budget was 
closed".to_string(),
+                        source: Some(Box::new(error)),
+                    }
+                })?;
+                #[cfg(test)]
+                let _query_io_probe_guard = match &self.query_io_probe {
+                    Some(probe) => Some(probe.enter().await),
+                    None => None,
+                };
+                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,
@@ -3277,6 +3357,72 @@ mod tests {
         assert_eq!(ranges, vec![RowRange::new(22, 26)]);
     }
 
+    #[tokio::test]
+    async fn test_multi_field_and_shares_query_concurrency_budget() {
+        let src = format!(
+            "{}/testdata/btree/btree_int_100_no_compress.bin",
+            env!("CARGO_MANIFEST_DIR")
+        );
+        let tmp = tempfile::tempdir().unwrap();
+        let index_dir = tmp.path().join("index");
+        std::fs::create_dir_all(&index_dir).unwrap();
+        let file_names: Vec<_> = (1..=4)
+            .map(|field_id| {
+                let file_name = format!("index_field{field_id}.bin");
+                std::fs::copy(&src, index_dir.join(&file_name)).unwrap();
+                file_name
+            })
+            .collect();
+
+        let table_path = format!("file://{}", tmp.path().display());
+        let file_io = crate::io::FileIOBuilder::new("file").build().unwrap();
+        let meta = BTreeIndexMeta::new(Some(le_int_key(0)), 
Some(le_int_key(198)), false);
+        let fields: Vec<_> = (0..4)
+            .map(|index| {
+                let field_id = index + 1;
+                DataField::new(
+                    field_id,
+                    format!("field{field_id}"),
+                    DataType::Int(crate::spec::IntType::new()),
+                )
+            })
+            .collect();
+        let entries: Vec<_> = file_names
+            .iter()
+            .enumerate()
+            .map(|(index, file_name)| {
+                make_global_index_entry(file_name, index as i32 + 1, 0, 99, 
&meta)
+            })
+            .collect();
+        let predicate = Predicate::and(
+            (0..4)
+                .map(|index| int_eq(&format!("field{}", index + 1), index, 50))
+                .collect(),
+        );
+
+        for (thread_num, expected_peak) in [(1, 1), (2, 2)] {
+            let mut scanner = GlobalIndexScanner::create(
+                &file_io,
+                &table_path,
+                thread_num,
+                i64::MAX,
+                i64::MAX,
+                &entries,
+                &fields,
+            )
+            .unwrap()
+            .unwrap();
+            let probe = Arc::new(QueryIoProbe::default());
+            scanner.query_io_probe = Some(Arc::clone(&probe));
+
+            let result = scanner.evaluate(&predicate).await.unwrap().unwrap();
+
+            assert_eq!(result.row_ranges, vec![RowRange::new(25, 25)]);
+            assert_eq!(result.evaluated_field_ids, HashSet::from([1, 2, 3, 
4]));
+            assert_eq!(probe.peak(), expected_peak);
+        }
+    }
+
     /// Regression for the Between+remaining bug in `evaluate_leaf`. When a
     /// native `Between` leaf is paired with another conjunct (e.g. `id >= 0`),
     /// and the file's b-tree key range falls **outside** the Between range

Reply via email to