Skip to content

[SPARK-59911][PYTHON] Migrate SQL_SCALAR_PANDAS_UDF to a pandas eval type handler - #59171

Closed
Yicong-Huang wants to merge 1 commit into
apache:masterfrom
Yicong-Huang:pandas-handler-base
Closed

Yicong-Huang wants to merge 1 commit into
apache:masterfrom
Yicong-Huang:pandas-handler-base

Conversation

@Yicong-Huang

@Yicong-Huang Yicong-Huang commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

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.

@Yicong-Huang
Yicong-Huang marked this pull request as draft September 30, 2026 22:51
@Yicong-Huang Yicong-Huang changed the title [SPARK-59911][PYTHON] Introduce a pandas eval-type handler base and migrate SQL_SCALAR_PANDAS_UDF [WIP][SPARK-59911][PYTHON] Introduce a pandas eval-type handler base and migrate SQL_SCALAR_PANDAS_UDF Sep 30, 2026
@Yicong-Huang
Yicong-Huang force-pushed the pandas-handler-base branch 3 times, most recently from ab7df7e to 83f56c0 Compare October 1, 2026 19:22
@Yicong-Huang Yicong-Huang changed the title [WIP][SPARK-59911][PYTHON] Introduce a pandas eval-type handler base and migrate SQL_SCALAR_PANDAS_UDF [WIP][SPARK-59911][PYTHON] Migrate SQL_SCALAR_PANDAS_UDF to a pandas eval type handler Oct 1, 2026
@Yicong-Huang
Yicong-Huang marked this pull request as ready for review October 1, 2026 19:46
@Yicong-Huang Yicong-Huang changed the title [WIP][SPARK-59911][PYTHON] Migrate SQL_SCALAR_PANDAS_UDF to a pandas eval type handler [SPARK-59911][PYTHON] Migrate SQL_SCALAR_PANDAS_UDF to a pandas eval type handler Oct 1, 2026
Co-authored-by: Isaac <no-reply@databricks.com>
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>
@Yicong-Huang

Copy link
Copy Markdown
Contributor Author

Merge Summary:

Posted by merge_spark_pr.py

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants