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()