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 b4a950bce2 [python][daft] Push filters into Paimon scan planning 
(#8744)
b4a950bce2 is described below

commit b4a950bce26f7d66e19bf7fa542eb29f958b6995
Author: QuakeWang <[email protected]>
AuthorDate: Mon Jul 20 21:35:16 2026 +0800

    [python][daft] Push filters into Paimon scan planning (#8744)
---
 paimon-python/pypaimon/daft/daft_datasource.py     | 45 +--------------------
 paimon-python/pypaimon/daft/daft_paimon.py         | 25 +++++-------
 .../pypaimon/tests/daft/daft_data_test.py          | 22 +++--------
 .../pypaimon/tests/daft/daft_integration_test.py   | 46 ++++++++++++++++++++++
 4 files changed, 61 insertions(+), 77 deletions(-)

diff --git a/paimon-python/pypaimon/daft/daft_datasource.py 
b/paimon-python/pypaimon/daft/daft_datasource.py
index f7232de0d1..d2175ac8e2 100644
--- a/paimon-python/pypaimon/daft/daft_datasource.py
+++ b/paimon-python/pypaimon/daft/daft_datasource.py
@@ -559,9 +559,6 @@ class PaimonDataSource(DataSource):
         table_path = getattr(table, "table_path", None)
         self._table_path = str(table_path) if table_path is not None else None
         self._table_options = _extract_table_options(table)
-        self._pushed_filters: list[PyExpr] | None = None
-        self._paimon_predicate: Predicate | None = None
-        self._remaining_filters: list[PyExpr] | None = None
         self._init_table(table)
 
     def __getstate__(self) -> dict[str, Any]:
@@ -573,9 +570,6 @@ class PaimonDataSource(DataSource):
             "_table_identifier": self._table_identifier,
             "_table_path": self._table_path,
             "_table_options": self._table_options,
-            "_pushed_filters": self._pushed_filters,
-            "_paimon_predicate": self._paimon_predicate,
-            "_remaining_filters": self._remaining_filters,
         }
 
     def __setstate__(self, state: dict[str, Any]) -> None:
@@ -585,9 +579,6 @@ class PaimonDataSource(DataSource):
         self._table_identifier = state["_table_identifier"]
         self._table_path = state["_table_path"]
         self._table_options = state["_table_options"]
-        self._pushed_filters = state.get("_pushed_filters")
-        self._paimon_predicate = state["_paimon_predicate"]
-        self._remaining_filters = state["_remaining_filters"]
 
         table = _load_table(
             self._table_catalog_options,
@@ -673,27 +664,6 @@ class PaimonDataSource(DataSource):
         partition_key_names = set(self._table.partition_keys)
         return [PartitionField.create(f) for f in self._schema if f.name in 
partition_key_names]
 
-    def push_filters(self, filters: list[PyExpr]) -> tuple[list[PyExpr], 
list[PyExpr]]:
-        """Push filters down to Paimon scan.
-
-        Converts Daft expressions to Paimon predicates where possible.
-        Returns (pushed_filters, remaining_filters).
-        """
-        pushed_filters, remaining_filters, paimon_predicate = 
convert_filters_to_paimon(self._table, filters)
-
-        self._pushed_filters = pushed_filters
-        self._paimon_predicate = paimon_predicate
-        self._remaining_filters = remaining_filters
-
-        if pushed_filters:
-            logger.debug(
-                "Paimon filter pushdown: %d filters pushed, %d remaining",
-                len(pushed_filters),
-                len(remaining_filters),
-            )
-
-        return pushed_filters, remaining_filters
-
     def _read_table_for_scan(self) -> FileStoreTable:
         if self._has_blob_columns:
             return self._table.copy({"blob-as-descriptor": "true"})
@@ -1081,12 +1051,6 @@ class PaimonDataSource(DataSource):
         return self._format_pyexprs([getattr(pushdowns.partition_filters, 
"_expr", pushdowns.partition_filters)])
 
     def _filter_pushdown_explain(self, pushdowns: Pushdowns) -> 
tuple[list[str], list[str]]:
-        if self._remaining_filters is not None:
-            return (
-                self._format_pyexprs(self._pushed_filters or []),
-                self._format_pyexprs(self._remaining_filters),
-            )
-
         if pushdowns.filters is None:
             return [], []
 
@@ -1246,8 +1210,6 @@ class PaimonDataSource(DataSource):
         )
 
     def _pushdown_filter_state(self, pushdowns: Pushdowns) -> tuple[Predicate 
| None, bool]:
-        if self._remaining_filters is not None:
-            return self._paimon_predicate, not self._remaining_filters
         if pushdowns.filters is None:
             return None, True
 
@@ -1260,9 +1222,7 @@ class PaimonDataSource(DataSource):
             return None
         if not self._can_plan_predicate(pushdown_predicate):
             return None
-        if self._paimon_predicate is not None or 
self._requires_fallback_reader():
-            return pushdown_predicate
-        return None
+        return pushdown_predicate
 
     def _partition_planning_predicate(self, pushdowns: Pushdowns) -> Predicate 
| None:
         """partition_filters -> Paimon predicate for plan-time pruning. 
Excludes
@@ -1307,9 +1267,6 @@ class PaimonDataSource(DataSource):
             return None
         return pushdowns.limit
 
-    def _requires_fallback_reader(self) -> bool:
-        return not self._is_parquet or self._has_blob_columns or 
self._table.is_primary_key_table
-
     def _can_plan_predicate(self, predicate: Predicate) -> bool:
         # Missing value null-count stats make isNull unsafe for scan planning.
         if not self._predicate_contains_is_null(predicate):
diff --git a/paimon-python/pypaimon/daft/daft_paimon.py 
b/paimon-python/pypaimon/daft/daft_paimon.py
index c4bfa7af17..e34b35b0d1 100644
--- a/paimon-python/pypaimon/daft/daft_paimon.py
+++ b/paimon-python/pypaimon/daft/daft_paimon.py
@@ -178,22 +178,19 @@ def _read_table(
     return _source_for_table(table, catalog_options=catalog_options, 
io_config=io_config).read()
 
 
-def _normalize_explain_filters(filters: Any) -> tuple[Any, list[Any]]:
+def _normalize_explain_filters(filters: Any) -> Any:
     if filters is None:
-        return None, []
+        return None
 
     if isinstance(filters, (list, tuple)):
         if not filters:
-            return None, []
-        filter_exprs = list(filters)
-        combined = filter_exprs[0]
-        for filter_expr in filter_exprs[1:]:
+            return None
+        combined = filters[0]
+        for filter_expr in filters[1:]:
             combined = combined & filter_expr
-    else:
-        filter_exprs = [filters]
-        combined = filters
+        return combined
 
-    return combined, [getattr(filter_expr, "_expr", filter_expr) for 
filter_expr in filter_exprs]
+    return filters
 
 
 def _explain_table(
@@ -216,14 +213,10 @@ def _explain_table(
         table, snapshot_id=snapshot_id, tag_name=tag_name, timestamp=timestamp
     )
     source = _source_for_table(table, catalog_options=catalog_options, 
io_config=io_config)
-    filter_expr, filter_pyexprs = _normalize_explain_filters(filters)
-    partition_filter_expr, _ = _normalize_explain_filters(partition_filters)
-    if filter_pyexprs:
-        source.push_filters(filter_pyexprs)
     return source.explain_scan(
         Pushdowns(
-            filters=filter_expr,
-            partition_filters=partition_filter_expr,
+            filters=_normalize_explain_filters(filters),
+            partition_filters=_normalize_explain_filters(partition_filters),
             columns=columns,
             limit=limit,
         ),
diff --git a/paimon-python/pypaimon/tests/daft/daft_data_test.py 
b/paimon-python/pypaimon/tests/daft/daft_data_test.py
index eeb0ecdc73..db18fce7da 100644
--- a/paimon-python/pypaimon/tests/daft/daft_data_test.py
+++ b/paimon-python/pypaimon/tests/daft/daft_data_test.py
@@ -106,7 +106,6 @@ async def _read_paimon_source_batches(
     filter_expr=None,
     columns=None,
     limit=None,
-    call_push_filters=True,
 ):
     from daft import context, runners
     from daft.daft import StorageConfig
@@ -118,11 +117,6 @@ async def _read_paimon_source_batches(
     storage_config = StorageConfig(runners.get_or_create_runner().name != 
"ray", io_config)
     source = PaimonDataSource(table, storage_config=storage_config, 
catalog_options={})
 
-    if filter_expr is not None and call_push_filters:
-        pushed_filters, remaining_filters = 
source.push_filters([filter_expr._expr])
-        assert pushed_filters
-        assert not remaining_filters
-
     pushdowns = Pushdowns(filters=filter_expr, columns=columns, limit=limit)
     return await _collect_paimon_source_batches(source, pushdowns)
 
@@ -328,8 +322,8 @@ def 
test_source_serialization_preserves_explicit_io_config_for_native_task(
         ) == explicit_values
 
 
-def 
test_read_paimon_source_serialization_preserves_pushed_filter_for_fallback(local_paimon_catalog):
-    """A serialized source must keep filters accepted by 
SupportsPushdownFilters."""
+def 
test_read_paimon_serialized_source_uses_current_filter_for_fallback(local_paimon_catalog):
+    """A serialized source must plan from the filters in the current 
request."""
     from daft import context, runners
     from daft.daft import StorageConfig
     from daft.io.pushdowns import Pushdowns
@@ -357,15 +351,12 @@ def 
test_read_paimon_source_serialization_preserves_pushed_filter_for_fallback(l
     io_config = context.get_context().daft_planning_config.default_io_config
     storage_config = StorageConfig(runners.get_or_create_runner().name != 
"ray", io_config)
     source = PaimonDataSource(table, storage_config=storage_config, 
catalog_options={})
-    pushed_filters, remaining_filters = source.push_filters([(col("id") == 
999)._expr])
-    assert pushed_filters
-    assert not remaining_filters
 
     restored = loads(dumps(source))
     batches = asyncio.run(
         _collect_paimon_source_batches(
             restored,
-            Pushdowns(filters=None, limit=1),
+            Pushdowns(filters=col("id") == 999, limit=1),
         )
     )
 
@@ -664,8 +655,8 @@ def 
test_read_paimon_pk_fallback_filters_before_projection(pk_table):
     assert batches == [{"name": ["new_a"], "id": [1]}]
 
 
-def 
test_read_paimon_fallback_plans_pushdown_filter_without_push_filters(local_paimon_catalog):
-    """Fallback planning must use Pushdowns.filters even if push_filters was 
not called."""
+def 
test_read_paimon_fallback_plans_current_pushdown_filter(local_paimon_catalog):
+    """Fallback planning must use the filter from the current request."""
     catalog, _ = local_paimon_catalog
     schema = pypaimon.Schema.from_pyarrow_schema(
         pa.schema([
@@ -688,7 +679,6 @@ def 
test_read_paimon_fallback_plans_pushdown_filter_without_push_filters(local_p
             table,
             filter_expr=col("id") == 999,
             limit=1,
-            call_push_filters=False,
         )
     )
 
@@ -716,7 +706,6 @@ def 
test_read_paimon_fallback_not_in_filter_excludes_nulls_before_limit(local_pa
             table,
             filter_expr=~col("id").is_in([1, 2]),
             limit=1,
-            call_push_filters=False,
         )
     )
 
@@ -906,7 +895,6 @@ def 
test_row_and_partition_filters_both_reach_planning_predicate(append_only_tab
     io_config = context.get_context().daft_planning_config.default_io_config
     storage_config = StorageConfig(runners.get_or_create_runner().name != 
"ray", io_config)
     source = PaimonDataSource(table, storage_config=storage_config, 
catalog_options={})
-    source.push_filters([(col("id") > 1)._expr])
     pushdowns = Pushdowns(filters=(col("id") > 1),
                           partition_filters=(col("dt") == "2024-01-02"))
 
diff --git a/paimon-python/pypaimon/tests/daft/daft_integration_test.py 
b/paimon-python/pypaimon/tests/daft/daft_integration_test.py
index de31cdfc49..dd28d5d0d6 100644
--- a/paimon-python/pypaimon/tests/daft/daft_integration_test.py
+++ b/paimon-python/pypaimon/tests/daft/daft_integration_test.py
@@ -572,6 +572,52 @@ def test_read_paimon_filter(catalog_options):
     }
 
 
+def test_read_paimon_filter_pruning_matches_explain(catalog_options, 
monkeypatch):
+    from daft.io.source import DataSourceTask
+
+    pa_schema = pa.schema([
+        ("id", pa.int64()),
+        ("name", pa.string()),
+    ])
+    identifier, table = _create_table(
+        catalog_options,
+        "read_filter_pruning",
+        pa_schema,
+        options={
+            "bucket": "1",
+            "file.format": "parquet",
+            "metadata.stats-mode": "full",
+        },
+    )
+    _write_arrow(table, pa.table({"id": [1], "name": ["file-a"]}, 
schema=pa_schema))
+    _write_arrow(table, pa.table({"id": [999], "name": ["file-b"]}, 
schema=pa_schema))
+
+    planned_paths = []
+    original_parquet = DataSourceTask.parquet
+
+    def capture_parquet_task(**kwargs):
+        planned_paths.append(kwargs["path"])
+        return original_parquet(**kwargs)
+
+    monkeypatch.setattr(DataSourceTask, "parquet", capture_parquet_task)
+    predicate = col("id") == 999
+
+    result = read_paimon(identifier, 
catalog_options).where(predicate).to_pydict()
+    explain = explain_paimon_scan(
+        identifier,
+        catalog_options,
+        filters=predicate,
+        verbose=True,
+    )
+
+    assert result == {"id": [999], "name": ["file-b"]}
+    assert len(planned_paths) == 1
+    assert explain.total_file_count == len(planned_paths)
+    assert explain.paimon_scan.file_skipping is not None
+    assert explain.paimon_scan.file_skipping.before == 2
+    assert explain.paimon_scan.file_skipping.after == 1
+
+
 def test_read_paimon_limit(catalog_options):
     data = pa.table(
         {

Reply via email to