[SPARK-59911][PYTHON] Migrate SQL_SCALAR_PANDAS_UDF to a pandas eval type handler - #59171
Closed
Yicong-Huang wants to merge 1 commit into
Closed
Yicong-Huang wants to merge 1 commit into
Yicong-Huang wants to merge 1 commit into
Conversation
Yicong-Huang
marked this pull request as draft
September 30, 2026 22:51
Yicong-Huang
force-pushed
the
pandas-handler-base
branch
3 times, most recently
from
October 1, 2026 19:22
ab7df7e to
83f56c0
Compare
Yicong-Huang
marked this pull request as ready for review
October 1, 2026 19:46
Yicong-Huang
force-pushed
the
pandas-handler-base
branch
from
October 2, 2026 00:47
83f56c0 to
5e5b34a
Compare
HyukjinKwon
approved these changes
Oct 2, 2026
Co-authored-by: Isaac <no-reply@databricks.com>
Yicong-Huang
force-pushed
the
pandas-handler-base
branch
from
October 2, 2026 05:27
5e5b34a to
7258747
Compare
Yicong-Huang
added a commit
that referenced
this pull request
Oct 2, 2026
…type handler ### What changes were proposed in this pull request? This is the first step of migrating the pandas UDF eval types onto the eval-handler framework (umbrella SPARK-59415), following the Arrow handler work already shipped in `pyspark/eval_handlers/_arrow.py`. New file `python/pyspark/eval_handlers/_pandas.py` adds `PandasScalarUDFHandler(BatchEvalTypeHandler)` for `SQL_SCALAR_PANDAS_UDF`, the pandas counterpart of `ArrowScalarUDFHandler`: convert each input RecordBatch to pandas Series (struct columns become DataFrames via `df_for_struct=True`), invoke each UDF once per batch, check the row count, and convert the pandas results back to one RecordBatch. It calls the Arrow<->pandas conversions in `pyspark.sql.conversion` (`ArrowToPandasConversion.to_pandas` / `PandasToArrowConversion.from_pandas`) directly with the runner_conf-derived parameters, exactly as `ArrowScalarUDFHandler` calls `ArrowBatchTransformer` directly -- no pandas-specific intermediate base and no conversion wrappers, so the handler reads as a straight pandas analog of the Arrow one. The class hierarchy after this PR (pandas mirrors Arrow -- both extend the category base directly): EvalTypeHandler BatchEvalTypeHandler ArrowScalarUDFHandler (SQL_SCALAR_ARROW_UDF, shipped) PandasScalarUDFHandler (SQL_SCALAR_PANDAS_UDF, new) The `SQL_SCALAR_PANDAS_UDF` branch is removed from `worker.py::read_udfs` (both the serializer-selection entry and the mapper block); `read_udfs` now dispatches it through `get_eval_type_handler`, exactly as the migrated Arrow types already do. All other pandas branches are left untouched. `read_single_udf` still handles `SQL_SCALAR_PANDAS_UDF` (the handler path calls it). One supporting fix in `pyspark/sql/conversion.py`: `ArrowToPandasConversion.to_pandas` declared `timezone: str`, but every sibling conversion method (and the `_convert_array` it delegates to) types it `Optional[str]`, and the worker passes `runner_conf.timezone` which is `Optional[str]`. The handler surfaces this (it is type-checked, unlike `worker.py`), so the annotation is corrected to `Optional[str]`. Annotation-only; no runtime change. This is a pure refactor with no behavior change: the handler issues the same `ArrowToPandasConversion.to_pandas` / `PandasToArrowConversion.from_pandas` calls with the same parameters, in the same validation order, raising the same error types -- it is the previous `read_udfs` mapper body moved into the handler. The one structural change -- inlining the args/kwargs offsets (`f(*pos, **kw)`) instead of pre-combining them through `wrap_kwargs_support` -- is equivalent (`wrap_kwargs_support` builds the same call) and matches how `ArrowScalarUDFHandler` already does it. The handler serializer is `ArrowStreamSerializer(write_start_stream=True)`, the same serializer the legacy path selected for this eval type. The handler passes `prefers_large_types=runner_conf.use_large_var_types` inline, matching the current mapper. Note the pre-existing inconsistencies in the other (untouched) pandas branches -- `GROUPED_AGG_PANDAS_ITER` pins `prefers_large_types=False`, and the stateful branches omit it (use the default). If these turn out to be bugs they should be fixed in a separate ticket, not here. #### Sketch: how SQL_GROUPED_MAP_PANDAS_UDF(201) would follow this pattern Confirming the pattern generalizes to the Grouped category. `SQL_GROUPED_MAP_PANDAS_UDF` is a grouped type (input `Iterator[Iterator[pa.RecordBatch]]`, `ArrowStreamGroupSerializer`), so it extends the grouped category base directly and calls the same conversions directly -- no pandas-specific base: class PandasGroupedMapUDFHandler(GroupedEvalTypeHandler["pa.RecordBatch"]): eval_type = PythonEvalType.SQL_GROUPED_MAP_PANDAS_UDF # __init__: parse key/value offsets (extract_key_value_indexes), build # output_schema = StructType([StructField("_0", return_type)]). # run(split_index, data): for each group, materialize its batches into one # pa.Table, then # all_series = ArrowToPandasConversion.to_pandas(table, timezone=rc.timezone, # prefer_int_ext_dtype=rc.prefer_int_ext_dtype) # build the value DataFrame from value_offsets, call the UDF (with the key # tuple when num_udf_args == 2), verify the pandas result, then # yield PandasToArrowConversion.from_pandas([result], output_schema, ...) The current `GROUPED_MAP_PANDAS_UDF` mapper already makes exactly those two calls, so moving it into a handler is the same mechanical move as this PR -- only the category base (`Grouped` vs `Batch`), the per-group loop, key handling, the pandas result verification, and the output-batch resizing wrapper differ. If a second pandas handler ends up repeating the same runner_conf-to-parameter argument list, that is the point to factor a shared helper (as the Arrow side keeps shared verifiers in `verification.py`) -- deferred here since there is a single caller. ### Why are the changes needed? Consolidating each Arrow/Pandas UDF eval type behind an `EvalTypeHandler` keeps the per-type logic together and shrinks the long `read_udfs` dispatch. This PR extends the pattern from the Arrow types to the first pandas type. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? - New unit tests `python/pyspark/eval_handlers/tests/test_pandas_eval_type_handlers.py` (built on a real `RunnerConf`): registration, per-batch invocation, output type coercion, multiple UDFs, keyword-offset binding, struct return (DataFrame), the three error paths (non-sized result, row-count mismatch, struct return that is not a DataFrame), and that the handler forwards `preferIntExtensionDtype` to the input conversion and `useLargeVarTypes` to the output conversion. - The existing per-eval-type suite `pyspark.sql.tests.pandas.test_pandas_udf_scalar` (the behavior gate for `SQL_SCALAR_PANDAS_UDF`) and the ASV per-eval-type benchmark (regression gate). ### Was this patch authored or co-authored using generative AI tooling? Yes, authored with Claude (Opus 4.8) via Claude Code. This pull request and its description were written by Isaac. Closes #59171 from Yicong-Huang/pandas-handler-base. Authored-by: Yicong Huang <17627829+Yicong-Huang@users.noreply.github.com> Signed-off-by: Yicong-Huang <17627829+Yicong-Huang@users.noreply.github.com> (cherry picked from commit 49a8389) Signed-off-by: Yicong-Huang <17627829+Yicong-Huang@users.noreply.github.com>
Contributor
Author
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
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?
This is the first step of migrating the pandas UDF eval types onto the eval-handler framework (umbrella SPARK-59415), following the Arrow handler work already shipped in
pyspark/eval_handlers/_arrow.py.New file
python/pyspark/eval_handlers/_pandas.pyaddsPandasScalarUDFHandler(BatchEvalTypeHandler)forSQL_SCALAR_PANDAS_UDF, the pandas counterpart ofArrowScalarUDFHandler: convert each input RecordBatch to pandas Series (struct columns become DataFrames viadf_for_struct=True), invoke each UDF once per batch, check the row count, and convert the pandas results back to one RecordBatch. It calls the Arrow<->pandas conversions inpyspark.sql.conversion(ArrowToPandasConversion.to_pandas/PandasToArrowConversion.from_pandas) directly with the runner_conf-derived parameters, exactly asArrowScalarUDFHandlercallsArrowBatchTransformerdirectly -- no pandas-specific intermediate base and no conversion wrappers, so the handler reads as a straight pandas analog of the Arrow one.The class hierarchy after this PR (pandas mirrors Arrow -- both extend the category base directly):
The
SQL_SCALAR_PANDAS_UDFbranch is removed fromworker.py::read_udfs(both the serializer-selection entry and the mapper block);read_udfsnow dispatches it throughget_eval_type_handler, exactly as the migrated Arrow types already do. All other pandas branches are left untouched.read_single_udfstill handlesSQL_SCALAR_PANDAS_UDF(the handler path calls it).One supporting fix in
pyspark/sql/conversion.py:ArrowToPandasConversion.to_pandasdeclaredtimezone: str, but every sibling conversion method (and the_convert_arrayit delegates to) types itOptional[str], and the worker passesrunner_conf.timezonewhich isOptional[str]. The handler surfaces this (it is type-checked, unlikeworker.py), so the annotation is corrected toOptional[str]. Annotation-only; no runtime change.This is a pure refactor with no behavior change: the handler issues the same
ArrowToPandasConversion.to_pandas/PandasToArrowConversion.from_pandascalls with the same parameters, in the same validation order, raising the same error types -- it is the previousread_udfsmapper body moved into the handler. The one structural change -- inlining the args/kwargs offsets (f(*pos, **kw)) instead of pre-combining them throughwrap_kwargs_support-- is equivalent (wrap_kwargs_supportbuilds the same call) and matches howArrowScalarUDFHandleralready does it. The handler serializer isArrowStreamSerializer(write_start_stream=True), the same serializer the legacy path selected for this eval type.The handler passes
prefers_large_types=runner_conf.use_large_var_typesinline, matching the current mapper. Note the pre-existing inconsistencies in the other (untouched) pandas branches --GROUPED_AGG_PANDAS_ITERpinsprefers_large_types=False, and the stateful branches omit it (use the default). If these turn out to be bugs they should be fixed in a separate ticket, not here.Sketch: how SQL_GROUPED_MAP_PANDAS_UDF(201) would follow this pattern
Confirming the pattern generalizes to the Grouped category.
SQL_GROUPED_MAP_PANDAS_UDFis a grouped type (inputIterator[Iterator[pa.RecordBatch]],ArrowStreamGroupSerializer), so it extends the grouped category base directly and calls the same conversions directly -- no pandas-specific base:The current
GROUPED_MAP_PANDAS_UDFmapper already makes exactly those two calls, so moving it into a handler is the same mechanical move as this PR -- only the category base (GroupedvsBatch), the per-group loop, key handling, the pandas result verification, and the output-batch resizing wrapper differ. If a second pandas handler ends up repeating the same runner_conf-to-parameter argument list, that is the point to factor a shared helper (as the Arrow side keeps shared verifiers inverification.py) -- deferred here since there is a single caller.Why are the changes needed?
Consolidating each Arrow/Pandas UDF eval type behind an
EvalTypeHandlerkeeps the per-type logic together and shrinks the longread_udfsdispatch. This PR extends the pattern from the Arrow types to the first pandas type.Does this PR introduce any user-facing change?
No.
How was this patch tested?
python/pyspark/eval_handlers/tests/test_pandas_eval_type_handlers.py(built on a realRunnerConf): registration, per-batch invocation, output type coercion, multiple UDFs, keyword-offset binding, struct return (DataFrame), the three error paths (non-sized result, row-count mismatch, struct return that is not a DataFrame), and that the handler forwardspreferIntExtensionDtypeto the input conversion anduseLargeVarTypesto the output conversion.pyspark.sql.tests.pandas.test_pandas_udf_scalar(the behavior gate forSQL_SCALAR_PANDAS_UDF) and the ASV per-eval-type benchmark (regression gate).Was this patch authored or co-authored using generative AI tooling?
Yes, authored with Claude (Opus 4.8) via Claude Code.
This pull request and its description were written by Isaac.