viirya opened a new pull request, #58978: URL: https://github.com/apache/spark/pull/58978
### What changes were proposed in this pull request? Add opt-in in-process Python UDF execution for [SPARK-59718](https://issues.apache.org/jira/browse/SPARK-59718). The new `pyspark.inprocess.inprocess_udf` API accepts one PyArrow array per input column and returns a PyArrow array. CPython runs inside the executor JVM through JEP, and Arrow C Data Interface exchanges compatible buffers without copying them across the JVM/Python boundary. Row-to-Arrow conversion and subsequent conversion back to Spark rows still incur costs. The implementation includes Catalyst planning and physical execution, an executor plugin for interpreter initialization and shutdown, result length/type validation, Arrow resource cleanup, and task cancellation checks. Integration and runtime tests join the existing `pyspark-sql` suite; the default Python 3.12 CI image installs JEP. No separate CI job is added. ASV benchmarks compare against worker-based Arrow UDFs, with a supplementary pandas UDF baseline. JEP is a `provided` dependency and is not bundled in the Spark distribution. Users must install and configure it explicitly. Existing Python UDF execution remains unchanged. The initial deployment supports one concurrent task per executor, typically configured with `spark.executor.cores == spark.task.cpus`. Application-level parallelism is retained through multiple executors. Interpreter pooling and Spark Connect support are outside this PR's scope. This is an explicit trade-off: additional executors introduce JVM overhead, and embedded Python provides weaker process isolation. A native Python crash can terminate the executor JVM. Arbitrary running Python/native code cannot be safely interrupted, so cancellation may wait for the current invocation to return. The feature targets users who can control executor configuration and accept these trade-offs. See the [design document](https://github.com/viirya/spark-1/blob/codex/inprocess-python-udf/docs/inprocess-python-udf-design.md) and [user guide](https://github.com/viirya/spark-1/blob/codex/inprocess-python-udf/docs/sql-pyspark-inprocess-udf.md). ### Why are the changes needed? Worker-based Arrow UDFs exchange data between executor JVMs and separate Python processes. For large inputs or outputs with relatively inexpensive computation, this transfer can contribute substantially to execution time. Preliminary benchmarks use matching Arrow operations on both sides, avoiding pandas conversion differences: * **Single task (`local[1]`):** the ASV benchmark measured approximately **3.0x speedup** for identity over one column of 1,000-character strings at 500,000, 1 million, and 2 million rows. Both paths use cached input and an output-consuming sink. * **Equal-resource parallel execution:** a separate experiment used four CPU slots, 4 GiB aggregate executor heap, a 1 GiB driver heap, and a 7 GiB total container memory limit. Each configuration supports four concurrent tasks. | Configuration | UDF | Executor layout | | --- | --- | --- | | A | Worker Arrow | One four-core executor | | B | Worker Arrow | Four single-core executors | | C | In-process | Four single-core executors | | Workload | A | B | C | C speedup over A | C speedup over B | | --- | ---: | ---: | ---: | ---: | ---: | | Identity over 800,000 strings of 1,000 characters | 0.336 s | 0.486 s | 0.261 s | 1.29x | 1.86x | | Sum of ten integer columns over 10 million rows | 1.382 s | 1.152 s | 0.570 s | 2.42x | 2.02x | These are medians of ten timed queries per case after warmup. The parallel experiment uses real separate executor JVMs via `local-cluster`, 16 partitions, and fully memory-cached inputs. All 960 measured tasks completed successfully without spilling, and output checksums were verified. These are local measurements, not GitHub Actions benchmark results or multi-host cluster results. The single-task and parallel experiments differ in environment, partitioning, and actual batch size; their ratios are not a controlled scaling comparison. The subsecond string queries show noticeable variability. Broader workload and distributed-cluster evaluation remains necessary. ### Does this PR introduce _any_ user-facing change? Yes. Users can explicitly define an in-process Arrow UDF: ```python import pyarrow.compute as pc from pyspark.inprocess import inprocess_udf from pyspark.sql.types import LongType @inprocess_udf(return_type=LongType()) def double(x): return pc.multiply(x, 2) spark.range(10).select(double("id")).show() ``` Deployment requires JEP and its native library on the executor load paths. Register `org.apache.spark.sql.execution.python.InProcessPythonPlugin` through `spark.plugins` before creating the SparkContext. The optional `spark.inprocess.python.sitePackages` setting supplies additional Python import paths. The result must match the declared type and input batch length. ### How was this patch tested? Added Scala suites covering planning/configuration, the Arrow bridge, and interpreter lifecycle/cancellation, plus Python integration and runtime contract suites. The standard Python runner passed **32 integration tests and 11 runtime tests, with no skips**, on both macOS/Python 3.13 and Linux/Python 3.12 with JEP 4.3.1: ```bash INPROCESS_TESTS=1 python/run-tests \ --python-executables="$(command -v python)" \ --testnames pyspark.sql.tests.test_inprocess_udf,pyspark.sql.tests.test_inprocess_runtime ``` Also checked optional dependency skipping, failure when JEP is missing but tests are explicitly required, Python formatting, test-module registration, and `git diff --check`. Ran the ASV Arrow-UDF comparison and the equal-resource multi-executor experiment summarized above. Measurements used the previously fully compiled Spark assembly; subsequent changes to benchmark/CI code and removal of legacy scripts did not change Scala implementation. The latest script-removal commit was checked for stale references and whitespace but did not trigger another benchmark run. Remote CI and GitHub Actions benchmarks have not yet been run for this PR. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude (original implementation; model/version not recorded) Generated-by: OpenAI Codex (review, implementation fixes, tests, CI integration, benchmarks, and documentation; model/version not recorded) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
