Repository navigation
Commit 49a8389
committed
[SPARK-59911][PYTHON] Migrate SQL_SCALAR_PANDAS_UDF to a pandas eval 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>1 parent 9ca6037 commit 49a8389
6 files changed
Lines changed: 293 additions & 78 deletions
File tree
- dev/sparktestsupport
- python/pyspark
- eval_handlers
- tests
- sql
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
665 | 665 | | |
666 | 666 | | |
667 | 667 | | |
| 668 | + | |
668 | 669 | | |
669 | 670 | | |
670 | 671 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
24 | 24 | | |
25 | 25 | | |
26 | 26 | | |
27 | | - | |
28 | | - | |
| 27 | + | |
| 28 | + | |
29 | 29 | | |
30 | 30 | | |
31 | | - | |
| 31 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
Lines changed: 162 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
| 125 | + | |
| 126 | + | |
| 127 | + | |
| 128 | + | |
| 129 | + | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
| 134 | + | |
| 135 | + | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
| 148 | + | |
| 149 | + | |
| 150 | + | |
| 151 | + | |
| 152 | + | |
| 153 | + | |
| 154 | + | |
| 155 | + | |
| 156 | + | |
| 157 | + | |
| 158 | + | |
| 159 | + | |
| 160 | + | |
| 161 | + | |
| 162 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1824 | 1824 | | |
1825 | 1825 | | |
1826 | 1826 | | |
1827 | | - | |
| 1827 | + | |
1828 | 1828 | | |
1829 | 1829 | | |
1830 | 1830 | | |
| |||
1838 | 1838 | | |
1839 | 1839 | | |
1840 | 1840 | | |
1841 | | - | |
1842 | | - | |
| 1841 | + | |
| 1842 | + | |
| 1843 | + | |
1843 | 1844 | | |
1844 | 1845 | | |
1845 | 1846 | | |
| |||
0 commit comments