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 6e39faa8 feat(file_index): add predicate evaluation foundation (#721)
6e39faa8 is described below
commit 6e39faa8cfb86ec621fe8839d4c39c7eda312c28
Author: QuakeWang <[email protected]>
AuthorDate: Thu Aug 20 10:55:15 2026 +0800
feat(file_index): add predicate evaluation foundation (#721)
---
.../paimon/src/file_index/file_index_predicate.rs | 405 +++++++++++++++++++++
.../file_index/{mod.rs => file_index_reader.rs} | 21 +-
crates/paimon/src/file_index/file_index_result.rs | 138 +++++++
crates/paimon/src/file_index/mod.rs | 8 +
4 files changed, 570 insertions(+), 2 deletions(-)
diff --git a/crates/paimon/src/file_index/file_index_predicate.rs
b/crates/paimon/src/file_index/file_index_predicate.rs
new file mode 100644
index 00000000..2d24f2a9
--- /dev/null
+++ b/crates/paimon/src/file_index/file_index_predicate.rs
@@ -0,0 +1,405 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use std::collections::HashMap;
+
+use crate::file_index::file_index_reader::FileIndexReader;
+use crate::file_index::file_index_result::FileIndexResult;
+use crate::spec::{DataType, Datum, Predicate, PredicateOperator};
+
+/// Evaluates predicate trees against file index readers grouped by column.
+pub(crate) struct FileIndexPredicate {
+ column_readers: HashMap<String, Vec<Box<dyn FileIndexReader>>>,
+}
+
+impl FileIndexPredicate {
+ /// Creates an evaluator from the index readers available for each column.
+ pub(crate) fn new(column_readers: HashMap<String, Vec<Box<dyn
FileIndexReader>>>) -> Self {
+ Self { column_readers }
+ }
+
+ /// Evaluates a predicate without reading data outside the supplied
indexes.
+ pub(crate) fn evaluate(&self, predicate: &Predicate) -> FileIndexResult {
+ match predicate {
+ Predicate::AlwaysTrue => FileIndexResult::Remain,
+ Predicate::AlwaysFalse => FileIndexResult::Skip,
+ Predicate::And(children) => self.evaluate_and(children),
+ Predicate::Or(children) => self.evaluate_or(children),
+ Predicate::Not(inner) => self.evaluate_not(inner),
+ Predicate::Leaf {
+ column,
+ index,
+ data_type,
+ op,
+ literals,
+ } => self.evaluate_leaf(column, *index, data_type, *op, literals),
+ }
+ }
+
+ fn evaluate_and(&self, predicates: &[Predicate]) -> FileIndexResult {
+ let mut result = FileIndexResult::Remain;
+ for predicate in predicates {
+ result = result.and(self.evaluate(predicate));
+ if !result.remain() {
+ break;
+ }
+ }
+ result
+ }
+
+ fn evaluate_or(&self, predicates: &[Predicate]) -> FileIndexResult {
+ let mut result = FileIndexResult::Skip;
+ for predicate in predicates {
+ result = result.or(self.evaluate(predicate));
+ if matches!(&result, FileIndexResult::Remain) {
+ break;
+ }
+ }
+ result
+ }
+
+ fn evaluate_not(&self, predicate: &Predicate) -> FileIndexResult {
+ match predicate {
+ Predicate::AlwaysTrue => FileIndexResult::Skip,
+ Predicate::AlwaysFalse => FileIndexResult::Remain,
+ Predicate::Not(inner) => self.evaluate(inner),
+ _ => FileIndexResult::Remain,
+ }
+ }
+
+ fn evaluate_leaf(
+ &self,
+ column: &str,
+ index: usize,
+ data_type: &DataType,
+ operator: PredicateOperator,
+ literals: &[Datum],
+ ) -> FileIndexResult {
+ let Some(readers) = self.column_readers.get(column) else {
+ return FileIndexResult::Remain;
+ };
+
+ let mut result = FileIndexResult::Remain;
+ for reader in readers {
+ result = result.and(reader.evaluate(column, index, data_type,
operator, literals));
+ if !result.remain() {
+ break;
+ }
+ }
+ result
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use std::sync::atomic::{AtomicUsize, Ordering};
+ use std::sync::Arc;
+
+ use roaring::RoaringBitmap;
+
+ use super::*;
+ use crate::spec::{DataType, Datum, IntType, PredicateOperator};
+
+ struct MockReader {
+ supported_operator: PredicateOperator,
+ result: FileIndexResult,
+ calls: Arc<AtomicUsize>,
+ }
+
+ impl FileIndexReader for MockReader {
+ fn evaluate(
+ &self,
+ _column: &str,
+ _index: usize,
+ _data_type: &DataType,
+ operator: PredicateOperator,
+ _literals: &[Datum],
+ ) -> FileIndexResult {
+ self.calls.fetch_add(1, Ordering::SeqCst);
+ if operator == self.supported_operator {
+ self.result.clone()
+ } else {
+ FileIndexResult::Remain
+ }
+ }
+ }
+
+ struct DefaultReader;
+
+ impl FileIndexReader for DefaultReader {}
+
+ struct AssertingReader;
+
+ impl FileIndexReader for AssertingReader {
+ fn evaluate(
+ &self,
+ column: &str,
+ index: usize,
+ data_type: &DataType,
+ operator: PredicateOperator,
+ literals: &[Datum],
+ ) -> FileIndexResult {
+ assert_eq!(column, "a");
+ assert_eq!(index, 7);
+ assert_eq!(data_type, &int_type());
+ assert_eq!(operator, PredicateOperator::Eq);
+ assert_eq!(literals, &[Datum::Int(42)]);
+ selection([1, 3])
+ }
+ }
+
+ fn int_type() -> DataType {
+ DataType::Int(IntType::new())
+ }
+
+ fn leaf(column: &str) -> Predicate {
+ leaf_with_operator(column, PredicateOperator::Eq)
+ }
+
+ fn leaf_with_operator(column: &str, operator: PredicateOperator) ->
Predicate {
+ Predicate::Leaf {
+ column: column.to_string(),
+ index: 7,
+ data_type: int_type(),
+ op: operator,
+ literals: vec![Datum::Int(42)],
+ }
+ }
+
+ fn selection(rows: impl IntoIterator<Item = u32>) -> FileIndexResult {
+ FileIndexResult::Selection(rows.into_iter().collect::<RoaringBitmap>())
+ }
+
+ fn mock_reader(result: FileIndexResult, calls: &Arc<AtomicUsize>) ->
Box<dyn FileIndexReader> {
+ Box::new(MockReader {
+ supported_operator: PredicateOperator::Eq,
+ result,
+ calls: Arc::clone(calls),
+ })
+ }
+
+ fn evaluator(
+ readers: impl IntoIterator<Item = (String, Vec<Box<dyn
FileIndexReader>>)>,
+ ) -> FileIndexPredicate {
+ FileIndexPredicate::new(readers.into_iter().collect())
+ }
+
+ #[test]
+ fn test_evaluate_constants_and_empty_compounds() {
+ let evaluator = FileIndexPredicate::new(HashMap::new());
+
+ assert_eq!(
+ evaluator.evaluate(&Predicate::AlwaysTrue),
+ FileIndexResult::Remain
+ );
+ assert_eq!(
+ evaluator.evaluate(&Predicate::AlwaysFalse),
+ FileIndexResult::Skip
+ );
+ assert_eq!(
+ evaluator.evaluate(&Predicate::And(vec![])),
+ FileIndexResult::Remain
+ );
+ assert_eq!(
+ evaluator.evaluate(&Predicate::Or(vec![])),
+ FileIndexResult::Skip
+ );
+ assert_eq!(
+ evaluator.evaluate(&Predicate::And(vec![
+ Predicate::AlwaysTrue,
+ Predicate::AlwaysFalse,
+ ])),
+ FileIndexResult::Skip
+ );
+ assert_eq!(
+ evaluator.evaluate(&Predicate::Or(vec![
+ Predicate::AlwaysFalse,
+ Predicate::AlwaysTrue,
+ ])),
+ FileIndexResult::Remain
+ );
+ }
+
+ #[test]
+ fn test_leaf_passes_existing_predicate_fields_to_reader() {
+ let evaluator = evaluator([(
+ "a".to_string(),
+ vec![Box::new(AssertingReader) as Box<dyn FileIndexReader>],
+ )]);
+
+ assert_eq!(evaluator.evaluate(&leaf("a")), selection([1, 3]));
+ }
+
+ #[test]
+ fn test_missing_reader_and_unsupported_operator_remain() {
+ let evaluator = evaluator([
+ (
+ "a".to_string(),
+ vec![Box::new(DefaultReader) as Box<dyn FileIndexReader>],
+ ),
+ ("empty".to_string(), vec![]),
+ ]);
+
+ assert_eq!(
+ evaluator.evaluate(&leaf("missing")),
+ FileIndexResult::Remain
+ );
+ assert_eq!(evaluator.evaluate(&leaf("empty")),
FileIndexResult::Remain);
+ assert_eq!(
+ evaluator.evaluate(&leaf_with_operator("a",
PredicateOperator::Gt)),
+ FileIndexResult::Remain
+ );
+ }
+
+ #[test]
+ fn test_leaf_intersects_reader_selections() {
+ let first_calls = Arc::new(AtomicUsize::new(0));
+ let second_calls = Arc::new(AtomicUsize::new(0));
+ let evaluator = evaluator([(
+ "a".to_string(),
+ vec![
+ mock_reader(selection([1, 2, 3]), &first_calls),
+ mock_reader(selection([2, 3, 4]), &second_calls),
+ ],
+ )]);
+
+ assert_eq!(evaluator.evaluate(&leaf("a")), selection([2, 3]));
+ assert_eq!(first_calls.load(Ordering::SeqCst), 1);
+ assert_eq!(second_calls.load(Ordering::SeqCst), 1);
+ }
+
+ #[test]
+ fn test_leaf_combines_readers_and_short_circuits() {
+ let first_calls = Arc::new(AtomicUsize::new(0));
+ let second_calls = Arc::new(AtomicUsize::new(0));
+ let evaluator = evaluator([(
+ "a".to_string(),
+ vec![
+ mock_reader(FileIndexResult::Skip, &first_calls),
+ mock_reader(FileIndexResult::Remain, &second_calls),
+ ],
+ )]);
+
+ assert_eq!(evaluator.evaluate(&leaf("a")), FileIndexResult::Skip);
+ assert_eq!(first_calls.load(Ordering::SeqCst), 1);
+ assert_eq!(second_calls.load(Ordering::SeqCst), 0);
+ }
+
+ #[test]
+ fn test_recursive_and_or_selection_combination() {
+ let a_calls = Arc::new(AtomicUsize::new(0));
+ let b_calls = Arc::new(AtomicUsize::new(0));
+ let c_calls = Arc::new(AtomicUsize::new(0));
+ let evaluator = evaluator([
+ (
+ "a".to_string(),
+ vec![mock_reader(selection([1, 2, 3]), &a_calls)],
+ ),
+ (
+ "b".to_string(),
+ vec![mock_reader(selection([2, 3, 4]), &b_calls)],
+ ),
+ (
+ "c".to_string(),
+ vec![mock_reader(selection([3, 4]), &c_calls)],
+ ),
+ ]);
+ let predicate = Predicate::And(vec![Predicate::Or(vec![leaf("a"),
leaf("b")]), leaf("c")]);
+
+ assert_eq!(evaluator.evaluate(&predicate), selection([3, 4]));
+ assert_eq!(a_calls.load(Ordering::SeqCst), 1);
+ assert_eq!(b_calls.load(Ordering::SeqCst), 1);
+ assert_eq!(c_calls.load(Ordering::SeqCst), 1);
+ }
+
+ #[test]
+ fn test_and_short_circuits_remaining_predicates() {
+ let first_calls = Arc::new(AtomicUsize::new(0));
+ let second_calls = Arc::new(AtomicUsize::new(0));
+ let evaluator = evaluator([
+ (
+ "a".to_string(),
+ vec![mock_reader(FileIndexResult::Skip, &first_calls)],
+ ),
+ (
+ "b".to_string(),
+ vec![mock_reader(FileIndexResult::Remain, &second_calls)],
+ ),
+ ]);
+
+ assert_eq!(
+ evaluator.evaluate(&Predicate::And(vec![leaf("a"), leaf("b")])),
+ FileIndexResult::Skip
+ );
+ assert_eq!(first_calls.load(Ordering::SeqCst), 1);
+ assert_eq!(second_calls.load(Ordering::SeqCst), 0);
+ }
+
+ #[test]
+ fn test_or_short_circuits_after_remain() {
+ let first_calls = Arc::new(AtomicUsize::new(0));
+ let second_calls = Arc::new(AtomicUsize::new(0));
+ let evaluator = evaluator([
+ (
+ "a".to_string(),
+ vec![mock_reader(FileIndexResult::Remain, &first_calls)],
+ ),
+ (
+ "b".to_string(),
+ vec![mock_reader(FileIndexResult::Skip, &second_calls)],
+ ),
+ ]);
+
+ assert_eq!(
+ evaluator.evaluate(&Predicate::Or(vec![leaf("a"), leaf("b")])),
+ FileIndexResult::Remain
+ );
+ assert_eq!(first_calls.load(Ordering::SeqCst), 1);
+ assert_eq!(second_calls.load(Ordering::SeqCst), 0);
+ }
+
+ #[test]
+ fn test_not_fails_open_except_for_safe_cases() {
+ let calls = Arc::new(AtomicUsize::new(0));
+ let evaluator = evaluator([(
+ "a".to_string(),
+ vec![mock_reader(FileIndexResult::Skip, &calls)],
+ )]);
+
+ assert_eq!(
+
evaluator.evaluate(&Predicate::Not(Box::new(Predicate::AlwaysTrue))),
+ FileIndexResult::Skip
+ );
+ assert_eq!(
+
evaluator.evaluate(&Predicate::Not(Box::new(Predicate::AlwaysFalse))),
+ FileIndexResult::Remain
+ );
+ assert_eq!(
+ evaluator.evaluate(&Predicate::Not(Box::new(leaf("a")))),
+ FileIndexResult::Remain
+ );
+ assert_eq!(calls.load(Ordering::SeqCst), 0);
+
+ assert_eq!(
+
evaluator.evaluate(&Predicate::Not(Box::new(Predicate::Not(Box::new(leaf(
+ "a"
+ )))))),
+ FileIndexResult::Skip
+ );
+ assert_eq!(calls.load(Ordering::SeqCst), 1);
+ }
+}
diff --git a/crates/paimon/src/file_index/mod.rs
b/crates/paimon/src/file_index/file_index_reader.rs
similarity index 56%
copy from crates/paimon/src/file_index/mod.rs
copy to crates/paimon/src/file_index/file_index_reader.rs
index ca9ee543..afc5c472 100644
--- a/crates/paimon/src/file_index/mod.rs
+++ b/crates/paimon/src/file_index/file_index_reader.rs
@@ -15,5 +15,22 @@
// specific language governing permissions and limitations
// under the License.
-mod file_index_format;
-pub use file_index_format::*;
+use crate::file_index::file_index_result::FileIndexResult;
+use crate::spec::{DataType, Datum, PredicateOperator};
+
+/// Evaluates leaf predicates against one concrete file index.
+pub(crate) trait FileIndexReader {
+ /// Evaluates the fields carried by [`crate::spec::Predicate::Leaf`].
+ ///
+ /// Readers must return [`FileIndexResult::Remain`] for unsupported
operators.
+ fn evaluate(
+ &self,
+ _column: &str,
+ _index: usize,
+ _data_type: &DataType,
+ _operator: PredicateOperator,
+ _literals: &[Datum],
+ ) -> FileIndexResult {
+ FileIndexResult::Remain
+ }
+}
diff --git a/crates/paimon/src/file_index/file_index_result.rs
b/crates/paimon/src/file_index/file_index_result.rs
new file mode 100644
index 00000000..0e8c8681
--- /dev/null
+++ b/crates/paimon/src/file_index/file_index_result.rs
@@ -0,0 +1,138 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use roaring::RoaringBitmap;
+
+/// Result of evaluating a predicate against file indexes.
+///
+/// Every result is a conservative candidate set and must contain every
matching
+/// row. `Remain` represents the full candidate set, `Skip` represents no rows,
+/// and `Selection` represents the rows which may match the predicate.
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub(crate) enum FileIndexResult {
+ /// The index cannot narrow the full candidate set.
+ Remain,
+ /// The file cannot contain matching rows.
+ Skip,
+ /// Only the listed zero-based row positions may match.
+ Selection(RoaringBitmap),
+}
+
+impl FileIndexResult {
+ /// Returns whether the file index result contains any possible matches.
+ pub(crate) fn remain(&self) -> bool {
+ match self {
+ Self::Remain => true,
+ Self::Skip => false,
+ Self::Selection(selection) => !selection.is_empty(),
+ }
+ }
+
+ /// Combines two file index results with logical AND.
+ pub(crate) fn and(self, other: Self) -> Self {
+ match (self, other) {
+ (Self::Skip, _) | (_, Self::Skip) => Self::Skip,
+ (Self::Remain, other) | (other, Self::Remain) => other,
+ (Self::Selection(mut left), Self::Selection(right)) => {
+ left &= right;
+ Self::Selection(left)
+ }
+ }
+ }
+
+ /// Combines two file index results with logical OR.
+ pub(crate) fn or(self, other: Self) -> Self {
+ match (self, other) {
+ (Self::Remain, _) | (_, Self::Remain) => Self::Remain,
+ (Self::Skip, other) | (other, Self::Skip) => other,
+ (Self::Selection(mut left), Self::Selection(right)) => {
+ left |= right;
+ Self::Selection(left)
+ }
+ }
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn selection(rows: impl IntoIterator<Item = u32>) -> FileIndexResult {
+ FileIndexResult::Selection(rows.into_iter().collect())
+ }
+
+ #[test]
+ fn test_remain() {
+ assert!(FileIndexResult::Remain.remain());
+ assert!(!FileIndexResult::Skip.remain());
+ assert!(selection([1]).remain());
+ assert!(!selection([]).remain());
+ }
+
+ #[test]
+ fn test_and() {
+ let rows = selection([1, 2]);
+
+ assert_eq!(
+ FileIndexResult::Remain.and(FileIndexResult::Remain),
+ FileIndexResult::Remain
+ );
+ assert_eq!(
+ FileIndexResult::Remain.and(FileIndexResult::Skip),
+ FileIndexResult::Skip
+ );
+ assert_eq!(FileIndexResult::Remain.and(rows.clone()), rows);
+ assert_eq!(rows.clone().and(FileIndexResult::Remain), rows);
+ assert_eq!(
+ FileIndexResult::Skip.and(rows.clone()),
+ FileIndexResult::Skip
+ );
+ assert_eq!(rows.and(FileIndexResult::Skip), FileIndexResult::Skip);
+ assert_eq!(
+ selection([1, 2, 3]).and(selection([2, 3, 4])),
+ selection([2, 3])
+ );
+ }
+
+ #[test]
+ fn test_or() {
+ let rows = selection([1, 2]);
+
+ assert_eq!(
+ FileIndexResult::Skip.or(FileIndexResult::Skip),
+ FileIndexResult::Skip
+ );
+ assert_eq!(
+ FileIndexResult::Skip.or(FileIndexResult::Remain),
+ FileIndexResult::Remain
+ );
+ assert_eq!(
+ FileIndexResult::Remain.or(rows.clone()),
+ FileIndexResult::Remain
+ );
+ assert_eq!(
+ rows.clone().or(FileIndexResult::Remain),
+ FileIndexResult::Remain
+ );
+ assert_eq!(FileIndexResult::Skip.or(rows.clone()), rows);
+ assert_eq!(rows.clone().or(FileIndexResult::Skip), rows);
+ assert_eq!(
+ selection([1, 2, 3]).or(selection([2, 3, 4])),
+ selection([1, 2, 3, 4])
+ );
+ }
+}
diff --git a/crates/paimon/src/file_index/mod.rs
b/crates/paimon/src/file_index/mod.rs
index ca9ee543..c50dfd6b 100644
--- a/crates/paimon/src/file_index/mod.rs
+++ b/crates/paimon/src/file_index/mod.rs
@@ -16,4 +16,12 @@
// under the License.
mod file_index_format;
+// Keep the predicate foundation crate-private until a concrete reader
validates its contract.
+#[allow(dead_code)]
+pub(crate) mod file_index_predicate;
+#[allow(dead_code)]
+pub(crate) mod file_index_reader;
+#[allow(dead_code)]
+pub(crate) mod file_index_result;
+
pub use file_index_format::*;