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::*;

Reply via email to