L. C. Hsieh created SPARK-59718:
-----------------------------------

             Summary: Add opt-in in-process Python UDF execution via Arrow C 
Data Interface
                 Key: SPARK-59718
                 URL: https://issues.apache.org/jira/browse/SPARK-59718
             Project: Spark
          Issue Type: Improvement
          Components: PySpark
    Affects Versions: 4.4.0
            Reporter: L. C. Hsieh
            Assignee: L. C. Hsieh



PySpark Arrow UDFs execute in separate Python worker processes and exchange 
Arrow data with executor JVMs. For workloads with large inputs or outputs and 
relatively inexpensive UDF computation, cross-process data transfer can 
contribute substantially to execution time.

This proposal adds an opt-in execution path that embeds CPython in the executor 
JVM using JEP and exchanges Arrow arrays through the Arrow C Data Interface. 
This allows sharing compatible Arrow buffers across the JVM/Python boundary 
without copying them, including buffers allocated by other libraries and 
wrapped by PyArrow. Row-to-Arrow conversion and any subsequent conversion back 
to Spark rows still incur costs.

**Proposed functionality**

- An `inprocess_udf` Python API accepting one PyArrow array per input column 
and returning a PyArrow array.
- Catalyst planning and physical execution support.
- An executor plugin managing interpreter initialization and shutdown.
- Validation of result lengths and types, Arrow resource cleanup, and task 
cancellation checks.
- Integration with the existing PySpark SQL test suite.
- Comparative benchmarks against worker-based Arrow UDFs.

Existing Python UDF execution remains unchanged. Users explicitly select the 
new API. JEP is a provided dependency, is not bundled with the Spark 
distribution, and must be installed and configured separately.

**Deployment model and trade-offs**

The initial implementation supports one concurrent task per executor. Users can 
retain application-level task parallelism by deploying multiple executors, 
typically with `spark.executor.cores` equal to `spark.task.cpus`. Interpreter 
pooling is outside the initial scope.

This configuration introduces trade-offs associated with additional JVM 
processes, including startup and memory overhead and changes to cache and 
shuffle behavior. The preliminary equal-resource experiments below demonstrate 
that this deployment model can retain a throughput advantage for the tested 
workloads.

Embedding Python also reduces process isolation: a native Python crash can 
terminate the executor JVM. Task cancellation checks are supported, but 
arbitrary running Python or native code cannot be safely interrupted; 
cancellation may need to wait for the current invocation to return.

The feature is intended for users who can control executor configuration and 
accept these isolation trade-offs in exchange for lower JVM/Python data 
exchange overhead. It is not intended to replace worker-based Arrow UDF 
execution universally.

The initial API targets classic PySpark. Spark Connect support is outside the 
initial scope.

**Preliminary performance results**

Benchmarks compare in-process execution with worker-based Arrow UDFs using the 
same Arrow inputs, outputs, and UDF operations, avoiding pandas conversion 
differences.

*Single-task comparison (`local[1]`)*

For an identity UDF over one column of 1,000-character strings, the PySpark ASV 
benchmark measured approximately **3.0× speedup** across datasets of 500,000, 1 
million, and 2 million rows.

Both implementations used the same single-task configuration, cached inputs, 
and an output-consuming sink.

*Equal-resource, multi-task comparison*

A separate experiment used four total CPU slots, 4 GiB aggregate executor heap, 
and a 7 GiB total container memory limit:

| Configuration | UDF execution | Executor layout |
|---|---|---|
| A | Worker-based Arrow UDF | One four-core executor |
| B | Worker-based Arrow UDF | Four single-core executors |
| C | In-process UDF | Four single-core executors |

All configurations supported four concurrent tasks. Median query times were:

| 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.29× | 1.86× |
| Sum of ten integer columns over 10 million rows | 1.382 s | 1.152 s | 0.570 s 
| 2.42× | 2.02× |

Results are medians of ten timed queries per configuration after warmup, using 
separate executor JVMs and 16 input partitions. All input partitions were 
cached in memory. Output checksums were validated, and all 960 measured tasks 
completed successfully without spilling.

The single-task and multi-task experiments differ in environment, partitioning, 
and actual batch sizes; their speedup ratios should not be interpreted as a 
controlled scaling comparison. The multi-task experiment ran on one machine in 
resource-limited containers, and the subsecond string queries showed noticeable 
variability. Broader workload and distributed-cluster evaluation remains 
necessary.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to