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]

Reply via email to