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 f0c8871 fix(scan): use evaluated fields for global index fallback
(#449)
f0c8871 is described below
commit f0c8871ea1440822c4b7225041b5939f0d6e8ca7
Author: QuakeWang <[email protected]>
AuthorDate: Sat Jul 4 23:26:40 2026 +0800
fix(scan): use evaluated fields for global index fallback (#449)
Track which predicate fields are actually evaluated by the global index
scanner and use only those fields when adding raw fallback ranges. This keeps
extra-field coverage from suppressing fallback rows before composite-key global
index reads are supported.
Add regression coverage for extra-field-only fallback and mixed
primary/extra predicates.
---
crates/paimon/src/table/global_index_scanner.rs | 165 ++++++++++++++++++++----
1 file changed, 142 insertions(+), 23 deletions(-)
diff --git a/crates/paimon/src/table/global_index_scanner.rs
b/crates/paimon/src/table/global_index_scanner.rs
index f604d7c..752ef81 100644
--- a/crates/paimon/src/table/global_index_scanner.rs
+++ b/crates/paimon/src/table/global_index_scanner.rs
@@ -37,8 +37,9 @@ use std::sync::Mutex;
type BoxedCmp = Box<dyn Fn(&[u8], &[u8]) -> Ordering + Send + Sync>;
-type EvaluateFuture<'a> =
- std::pin::Pin<Box<dyn std::future::Future<Output =
Result<Option<Vec<RowRange>>>> + Send + 'a>>;
+type EvaluateFuture<'a> = std::pin::Pin<
+ Box<dyn std::future::Future<Output =
Result<Option<GlobalIndexScanResult>>> + Send + 'a>,
+>;
type PredicateTuple<'a> = (PredicateOperator, &'a [Datum], &'a DataType);
@@ -46,6 +47,11 @@ const BTREE_INDEX_TYPE: &str = "btree";
const DELETION_VECTORS_INDEX_TYPE: &str = "DELETION_VECTORS";
const INDEX_DIR: &str = "index";
+struct GlobalIndexScanResult {
+ row_ranges: Vec<RowRange>,
+ evaluated_field_ids: HashSet<i32>,
+}
+
/// Evaluates global index predicates and returns matching row ranges.
///
/// The scanner filters index manifest entries for global index files,
@@ -145,7 +151,7 @@ impl GlobalIndexScanner {
/// Evaluate a predicate against the global indexes and return matching
row ranges.
/// Returns `None` if the predicate cannot be evaluated by the global
index.
- pub(crate) fn evaluate<'a>(&'a self, predicate: &'a Predicate) ->
EvaluateFuture<'a> {
+ fn evaluate<'a>(&'a self, predicate: &'a Predicate) -> EvaluateFuture<'a> {
Box::pin(async move {
match predicate {
Predicate::Leaf {
@@ -169,6 +175,12 @@ impl GlobalIndexScanner {
};
self.evaluate_leaf(entries, &[(*op, literals.as_slice(),
data_type)])
.await
+ .map(|ranges| {
+ ranges.map(|row_ranges| GlobalIndexScanResult {
+ row_ranges,
+ evaluated_field_ids: HashSet::from([field_id]),
+ })
+ })
}
Predicate::And(children) => {
// Group leaf predicates by field_id to reuse readers
@@ -201,45 +213,61 @@ impl GlobalIndexScanner {
non_leaf_children.push(child);
}
- let mut result: Option<Vec<RowRange>> = None;
+ let mut row_ranges: Option<Vec<RowRange>> = None;
+ let mut evaluated_field_ids = HashSet::new();
// Evaluate grouped leaves (one open per file)
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? {
- result = Some(match result {
+ row_ranges = Some(match row_ranges {
None => ranges,
Some(existing) =>
intersect_sorted_ranges(&existing, &ranges),
});
+ evaluated_field_ids.insert(*field_id);
}
}
}
// Evaluate non-leaf children recursively
for child in non_leaf_children {
- if let Some(ranges) = self.evaluate(child).await? {
- result = Some(match result {
- None => ranges,
- Some(existing) =>
intersect_sorted_ranges(&existing, &ranges),
+ if let Some(child_result) =
self.evaluate(child).await? {
+ row_ranges = Some(match row_ranges {
+ None => child_result.row_ranges,
+ Some(existing) => {
+ intersect_sorted_ranges(&existing,
&child_result.row_ranges)
+ }
});
+
evaluated_field_ids.extend(child_result.evaluated_field_ids);
}
}
- Ok(result)
+ Ok(row_ranges.map(|row_ranges| GlobalIndexScanResult {
+ row_ranges,
+ evaluated_field_ids,
+ }))
}
Predicate::Or(children) => {
let mut all_ranges: Vec<RowRange> = Vec::new();
+ let mut evaluated_field_ids = HashSet::new();
for child in children {
match self.evaluate(child).await? {
- Some(ranges) => all_ranges.extend(ranges),
+ Some(child_result) => {
+ all_ranges.extend(child_result.row_ranges);
+
evaluated_field_ids.extend(child_result.evaluated_field_ids);
+ }
None => return Ok(None),
}
}
- if all_ranges.is_empty() {
- Ok(Some(Vec::new()))
+ let row_ranges = if all_ranges.is_empty() {
+ Vec::new()
} else {
- Ok(Some(super::merge_row_ranges(all_ranges)))
- }
+ super::merge_row_ranges(all_ranges)
+ };
+ Ok(Some(GlobalIndexScanResult {
+ row_ranges,
+ evaluated_field_ids,
+ }))
}
_ => Ok(None),
}
@@ -425,6 +453,7 @@ impl GlobalIndexScanner {
/// data-file row ranges collected by the scan. The caller unions these
/// ranges with indexed matches, and the normal read filter evaluates the
/// predicate on the raw rows.
+ #[cfg(test)]
fn unindexed_ranges(
&self,
predicate: &Predicate,
@@ -433,21 +462,33 @@ impl GlobalIndexScanner {
data_ranges: &[RowRange],
) -> Result<Vec<RowRange>> {
let field_ids = self.collect_field_ids(predicate)?;
- Ok(unindexed_ranges_for_coverage(
+ Ok(self.unindexed_ranges_for_field_ids(&field_ids, search_mode,
next_row_id, data_ranges))
+ }
+
+ fn unindexed_ranges_for_field_ids(
+ &self,
+ field_ids: &HashSet<i32>,
+ search_mode: GlobalIndexSearchMode,
+ next_row_id: Option<i64>,
+ data_ranges: &[RowRange],
+ ) -> Vec<RowRange> {
+ unindexed_ranges_for_coverage(
&self.coverage_by_field,
- &field_ids,
+ field_ids,
search_mode,
next_row_id,
data_ranges,
- ))
+ )
}
+ #[cfg(test)]
fn collect_field_ids(&self, predicate: &Predicate) -> Result<HashSet<i32>>
{
let mut field_ids = HashSet::new();
self.collect_field_ids_inner(predicate, &mut field_ids)?;
Ok(field_ids)
}
+ #[cfg(test)]
fn collect_field_ids_inner(
&self,
predicate: &Predicate,
@@ -916,16 +957,17 @@ pub(crate) async fn evaluate_global_index(
let combined = Predicate::and(evaluation.predicates.to_vec());
- let mut row_ranges = match scanner.evaluate(&combined).await? {
- Some(row_ranges) => row_ranges,
+ let scan_result = match scanner.evaluate(&combined).await? {
+ Some(scan_result) => scan_result,
None => return Ok(None),
};
- row_ranges.extend(scanner.unindexed_ranges(
- &combined,
+ let mut row_ranges = scan_result.row_ranges;
+ row_ranges.extend(scanner.unindexed_ranges_for_field_ids(
+ &scan_result.evaluated_field_ids,
evaluation.search_mode,
evaluation.next_row_id,
evaluation.data_ranges,
- )?);
+ ));
Ok(Some(super::merge_row_ranges(row_ranges)))
}
@@ -1294,6 +1336,31 @@ mod tests {
assert!(ranges.is_empty());
}
+ #[tokio::test]
+ async fn
test_evaluate_extra_field_only_without_composite_reader_falls_back() {
+ let (file_io, table_path, file_name, _tmp) =
+ setup_testdata_table("btree_int_100_no_compress.bin");
+ let meta = BTreeIndexMeta::new(Some(le_int_key(0)),
Some(le_int_key(198)), false);
+ let mut entry = make_global_index_entry(&file_name, 1, 0, 99, &meta);
+ entry
+ .index_file
+ .global_index_meta
+ .as_mut()
+ .unwrap()
+ .extra_field_ids = Some(vec![2]);
+ let fields = two_field_schema_fields();
+ let predicates = vec![int_eq("value", 1, 50)];
+
+ let result =
+ evaluate_global_index_fast(&file_io, &table_path, &[entry],
&predicates, &fields)
+ .await
+ .unwrap();
+ assert!(
+ result.is_none(),
+ "extra-field-only predicates must fall back until composite-key
btree reads are supported"
+ );
+ }
+
#[tokio::test]
async fn test_evaluate_global_index_eq() {
let (file_io, table_path, file_name, _tmp) =
@@ -1347,6 +1414,58 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn
test_evaluate_global_index_and_uses_evaluated_field_coverage_for_raw_fallback()
{
+ 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();
+ std::fs::copy(&src, index_dir.join("index_part1.bin")).unwrap();
+ std::fs::copy(&src, index_dir.join("index_part2.bin")).unwrap();
+
+ 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 mut first = make_global_index_entry("index_part1.bin", 1, 0, 49,
&meta);
+ first
+ .index_file
+ .global_index_meta
+ .as_mut()
+ .unwrap()
+ .extra_field_ids = Some(vec![2]);
+ let second = make_global_index_entry("index_part2.bin", 1, 50, 99,
&meta);
+ let entries = vec![first, second];
+ let fields = two_field_schema_fields();
+
+ let predicates = vec![Predicate::and(vec![
+ int_eq("id", 0, 50),
+ int_eq("value", 1, 8),
+ ])];
+ let result = super::evaluate_global_index(super::GlobalIndexEvaluation
{
+ file_io: &file_io,
+ table_path: &table_path,
+ index_entries: &entries,
+ predicates: &predicates,
+ schema_fields: &fields,
+ search_mode: GlobalIndexSearchMode::Full,
+ next_row_id: Some(100),
+ data_ranges: &[],
+ })
+ .await
+ .unwrap();
+
+ assert_eq!(
+ result.unwrap(),
+ vec![RowRange::new(25, 25), RowRange::new(75, 75)],
+ "raw fallback should use only the id field that was actually
evaluated; \
+ the unevaluated extra field must not widen or narrow fallback
coverage"
+ );
+ }
+
#[tokio::test]
async fn test_evaluate_global_index_detail_mode_uses_data_ranges() {
let (file_io, table_path, file_name, _tmp) =