[SPARK-51271][PYTHON] Add filter pushdown API to Python Data Sources #49961
+777
−79
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this pull request?
Suggested reivew order (from high level to details)
datasource.py
: add filter pushdown to Python Data Source APItest_python_datasource.py
: tests for filter pushdownPythonScanBuilder.scala
: implement filter pushdown API in ScalaUserDefinedPythonDataSource.scala
(UserDefinedPythonDataSourceFilterPushdownRunner
),data_source_pushdown_filters.py
: communication between Python and Scalasql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/python
and inpython/pyspark/sql
: Changes to the sequence of interactions between Python and Scala to accommodate filter pushdownChanges to interactions between Python and Scala
Original sequence:
data:image/s3,"s3://crabby-images/c326c/c326c32df050dfc57891d60d606a6bf5db2b2cc3" alt="pyds old-2025-02-12-001557"
Updated sequence (new interactions are highlighted in yellow):
data source -> (partitions, read function)
transformation (plan_data_source_read.py
) into two steps:data source -> reader
(data_source_get_reader.py
) andreader -> (partitions, read function)
(plan_data_source_read.py
).Why are the changes needed?
Filter pushdown allows reducing the amount of data produced by the reader, by filtering rows directly in the data source scan. The reduction in the amount of data can improve query performance. This PR implements filter pushdown for Python Data Sources API using the existing Scala DS filter pushdown API. An upcoming PR will implement the actual filter types and the serialization of filters.
Does this PR introduce any user-facing change?
Yes. New API are added. See
datasource.py
for details.How was this patch tested?
Tests added to
test_python_datasource.py
.Was this patch authored or co-authored using generative AI tooling?
No