This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 46eae86dc4 [python] Apply filters when reading system tables (#10164)
46eae86dc4 is described below

commit 46eae86dc47ef7c92a971bfbedc622c8225d25c3
Author: jackylee <[email protected]>
AuthorDate: Fri Sep 25 21:10:13 2026 +0800

    [python] Apply filters when reading system tables (#10164)
---
 .../pypaimon/table/system/system_table_read.py     |  70 ++++++++++-
 .../pypaimon/tests/system/system_table_test.py     | 135 ++++++++++++++++++++-
 2 files changed, 200 insertions(+), 5 deletions(-)

diff --git a/paimon-python/pypaimon/table/system/system_table_read.py 
b/paimon-python/pypaimon/table/system/system_table_read.py
index c02258ac3f..e33717227e 100644
--- a/paimon-python/pypaimon/table/system/system_table_read.py
+++ b/paimon-python/pypaimon/table/system/system_table_read.py
@@ -32,7 +32,8 @@ import pyarrow
 
 from pypaimon.common.predicate import Predicate
 from pypaimon.read.split import Split
-from pypaimon.schema.data_types import DataField, PyarrowFieldParser
+from pypaimon.schema.data_types import (
+    AtomicType, DataField, PyarrowFieldParser)
 from pypaimon.table.system.system_table_scan import SystemSplit
 
 if TYPE_CHECKING:  # pragma: no cover - type-only import
@@ -40,10 +41,16 @@ if TYPE_CHECKING:  # pragma: no cover - type-only import
 
 
 _PREDICATE_NOT_SUPPORTED = (
-    "predicate pushdown on system tables is not supported yet; "
-    "filter the resulting PyArrow / pandas table on the client side"
+    "string-match predicate pushdown (starts_with / ends_with / contains / "
+    "like) on system tables is not supported yet; filter the resulting "
+    "PyArrow / pandas table on the client side"
 )
 
+# Value comparisons (equal / greater-than / is-in / ...) have no PyArrow
+# compute kernel for nested columns such as ``$files.write_cols`` or
+# ``$table_indexes.dv_ranges``; only null checks work there.
+_COMPARISON_NEEDS_SCALAR = ("isNull", "isNotNull")
+
 
 class SystemTableRead:
     """In-memory read pipeline backed by a single SystemSplit."""
@@ -95,12 +102,29 @@ class SystemTableRead:
         return con
 
     def _materialise(self, splits: List[Split]) -> pyarrow.Table:
-        if self.predicate is not None:
+        from pypaimon.read.push_down_utils import 
predicate_supports_arrow_filter
+
+        predicate = self.predicate
+        if predicate is not None and not isinstance(predicate, Predicate):
+            # Preserve the public error contract for non-Predicate filter
+            # objects (e.g. ``with_filter(object())``): predicate_supports_
+            # arrow_filter dereferences ``.method`` and would otherwise leak
+            # an internal AttributeError.
+            raise NotImplementedError(_PREDICATE_NOT_SUPPORTED)
+        if predicate is not None and not predicate_supports_arrow_filter(
+            predicate
+        ):
             raise NotImplementedError(_PREDICATE_NOT_SUPPORTED)
+        if predicate is not None:
+            self._reject_non_scalar_filter_fields(predicate)
 
         if not splits:
             return self._empty_table()
 
+        arrow_predicate = (
+            self.predicate.to_arrow() if self.predicate is not None else None
+        )
+
         projected_names = [f.name for f in self.read_type]
         slices: List[pyarrow.Table] = []
         for split in splits:
@@ -110,6 +134,10 @@ class SystemTableRead:
                     + type(split).__name__
                 )
             arrow_table = split.arrow_table()
+            # Filter on the full schema before projection so a predicate may
+            # reference columns that are not part of the projection.
+            if arrow_predicate is not None:
+                arrow_table = arrow_table.filter(arrow_predicate)
             if projected_names:
                 arrow_table = arrow_table.select(projected_names)
             slices.append(arrow_table)
@@ -119,6 +147,40 @@ class SystemTableRead:
             combined = combined.slice(0, self.limit)
         return combined
 
+    def _reject_non_scalar_filter_fields(self, predicate: Predicate) -> None:
+        """Reject value comparisons on nested system-table columns.
+
+        System tables expose list/map/row columns (e.g. ``$files.write_cols``,
+        ``$table_indexes.dv_ranges``). PyArrow has no equality / comparison /
+        is-in kernel for those types, so ``arrow_table.filter`` would raise an
+        internal ``ArrowNotImplementedError``. Surface the documented
+        ``NotImplementedError`` instead. Null checks are left alone because the
+        ``is_null`` / ``is_valid`` kernels do accept nested input.
+        """
+        field_types = {
+            f.name: f.type for f in self.system_table.row_type().fields
+        }
+        pending: List[Predicate] = [predicate]
+        while pending:
+            current = pending.pop()
+            if current.method in ("and", "or"):
+                pending.extend(current.literals or [])
+                continue
+            if current.method in _COMPARISON_NEEDS_SCALAR:
+                continue
+            field_type = field_types.get(current.field)
+            if field_type is not None and not isinstance(field_type, 
AtomicType):
+                raise NotImplementedError(
+                    "predicate pushdown of "
+                    + repr(current.method)
+                    + " on the non-scalar system-table column "
+                    + repr(current.field)
+                    + " (type "
+                    + type(field_type).__name__
+                    + ") is not supported; filter the resulting PyArrow / "
+                    + "pandas table on the client side"
+                )
+
     def _empty_table(self) -> pyarrow.Table:
         target_schema = PyarrowFieldParser.from_paimon_schema(self.read_type)
         return pyarrow.Table.from_arrays(
diff --git a/paimon-python/pypaimon/tests/system/system_table_test.py 
b/paimon-python/pypaimon/tests/system/system_table_test.py
index 60a505679a..2bf76319ef 100644
--- a/paimon-python/pypaimon/tests/system/system_table_test.py
+++ b/paimon-python/pypaimon/tests/system/system_table_test.py
@@ -21,7 +21,9 @@ import types
 import unittest
 
 from pypaimon.common.identifier import Identifier
-from pypaimon.schema.data_types import AtomicType, DataField, RowType
+from pypaimon.common.predicate_builder import PredicateBuilder
+from pypaimon.schema.data_types import (
+    ArrayType, AtomicType, DataField, RowType)
 from pypaimon.table.system.system_table import SystemTable
 
 
@@ -50,6 +52,53 @@ class _DummySystemTable(SystemTable):
         return pa.table({"key": [], "value": []})
 
 
+_DATA_ROW_TYPE = RowType(False, [
+    DataField(0, "id", AtomicType("INT", nullable=False)),
+    DataField(1, "name", AtomicType("STRING", nullable=True)),
+])
+
+
+class _DataSystemTable(SystemTable):
+    """A SystemTable with a few real rows, used to exercise read filtering."""
+
+    def system_table_name(self) -> str:
+        return "data"
+
+    def row_type(self) -> RowType:
+        return _DATA_ROW_TYPE
+
+    def _build_arrow_table(self):
+        import pyarrow as pa
+        return pa.table({
+            "id": pa.array([1, 2, 3], pa.int32()),
+            "name": pa.array(["a", "b", "c"]),
+        })
+
+
+_ARRAY_ROW_TYPE = RowType(False, [
+    DataField(0, "id", AtomicType("INT", nullable=False)),
+    DataField(1, "tags",
+              ArrayType(True, AtomicType("STRING", nullable=True))),
+])
+
+
+class _ArraySystemTable(SystemTable):
+    """A SystemTable with a list-typed column (mirrors $files.write_cols)."""
+
+    def system_table_name(self) -> str:
+        return "arr"
+
+    def row_type(self) -> RowType:
+        return _ARRAY_ROW_TYPE
+
+    def _build_arrow_table(self):
+        import pyarrow as pa
+        return pa.table({
+            "id": pa.array([1, 2], pa.int32()),
+            "tags": pa.array([["a"], ["b"]], pa.list_(pa.string())),
+        })
+
+
 def _fake_base(database: str = "db", table: str = "t", branch=None):
     """Construct a minimal stand-in for FileStoreTable.
 
@@ -109,5 +158,89 @@ class SystemTableTest(unittest.TestCase):
                           "method {}: {}".format(method_name, ctx.exception))
 
 
+class SystemTableFilterTest(unittest.TestCase):
+    """``with_filter`` on a system table applies the predicate at read time."""
+
+    def _predicate_builder(self) -> PredicateBuilder:
+        return PredicateBuilder(_DATA_ROW_TYPE.fields)
+
+    def _read(self, predicate, projection=None):
+        sys_table = _DataSystemTable(_fake_base())
+        rb = sys_table.new_read_builder()
+        if projection is not None:
+            rb = rb.with_projection(projection)
+        rb = rb.with_filter(predicate)
+        splits = rb.new_scan().plan().splits()
+        return rb.new_read().to_arrow(splits)
+
+    def test_equal_filter_selects_matching_rows(self):
+        table = self._read(self._predicate_builder().equal("id", 2))
+        self.assertEqual(table.column("id").to_pylist(), [2])
+        self.assertEqual(table.column("name").to_pylist(), ["b"])
+
+    def test_greater_than_filter(self):
+        table = self._read(self._predicate_builder().greater_than("id", 1))
+        self.assertEqual(sorted(table.column("id").to_pylist()), [2, 3])
+
+    def test_is_in_filter(self):
+        table = self._read(self._predicate_builder().is_in("id", [1, 3]))
+        self.assertEqual(sorted(table.column("id").to_pylist()), [1, 3])
+
+    def test_filter_column_may_be_absent_from_projection(self):
+        # Filter on `id` while projecting only `name`: the predicate is applied
+        # before projection, so this must not raise a missing-column error.
+        table = self._read(self._predicate_builder().equal("id", 3),
+                           projection=["name"])
+        self.assertEqual(table.column_names, ["name"])
+        self.assertEqual(table.column("name").to_pylist(), ["c"])
+
+    def test_string_match_filter_still_raises(self):
+        # starts_with / ends_with / contains / like are not safe as final
+        # Arrow row filters, so they still surface a clear NotImplementedError.
+        pred = self._predicate_builder().contains("name", "b")
+        sys_table = _DataSystemTable(_fake_base())
+        rb = sys_table.new_read_builder().with_filter(pred)
+        splits = rb.new_scan().plan().splits()
+        read = rb.new_read()
+        with self.assertRaises(NotImplementedError):
+            read.to_arrow(splits)
+
+    def test_comparison_on_array_column_raises(self):
+        # PyArrow's ArrowNotImplementedError (a NotImplementedError subclass)
+        # for a list column carries only a cryptic "no kernel" message. The
+        # read must reject the comparison up front with the documented,
+        # column-named message instead.
+        pred = PredicateBuilder(_ARRAY_ROW_TYPE.fields).equal("tags", ["a"])
+        sys_table = _ArraySystemTable(_fake_base())
+        rb = sys_table.new_read_builder().with_filter(pred)
+        splits = rb.new_scan().plan().splits()
+        read = rb.new_read()
+        with self.assertRaises(NotImplementedError) as ctx:
+            read.to_arrow(splits)
+        message = str(ctx.exception)
+        self.assertIn("non-scalar", message)
+        self.assertIn("tags", message)
+
+    def test_null_check_on_array_column_is_allowed(self):
+        # is_null / is_valid kernels accept nested input, so a null check on a
+        # list column is filtered normally rather than rejected.
+        pred = PredicateBuilder(_ARRAY_ROW_TYPE.fields).is_not_null("tags")
+        sys_table = _ArraySystemTable(_fake_base())
+        rb = sys_table.new_read_builder().with_filter(pred)
+        splits = rb.new_scan().plan().splits()
+        table = rb.new_read().to_arrow(splits)
+        self.assertEqual(sorted(table.column("id").to_pylist()), [1, 2])
+
+    def test_non_predicate_filter_raises(self):
+        # with_filter(object()) must keep the public NotImplementedError
+        # contract instead of leaking an internal AttributeError.
+        sys_table = _DataSystemTable(_fake_base())
+        rb = sys_table.new_read_builder().with_filter(object())
+        splits = rb.new_scan().plan().splits()
+        read = rb.new_read()
+        with self.assertRaises(NotImplementedError):
+            read.to_arrow(splits)
+
+
 if __name__ == "__main__":
     unittest.main()

Reply via email to