Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions .github/actions/rust-test/action.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,28 @@ runs:
steps:
# Note: cargo fmt check is now handled by the lint job that gates this workflow

- name: Set up embedded Python for native UDF tests
shell: bash
run: |
apt-get update
apt-get install -y --no-install-recommends python3 python3-dev python3-venv
python3 -m venv /tmp/comet-rust-python
/tmp/comet-rust-python/bin/pip install "pyarrow>=14" cloudpickle
echo "PYO3_PYTHON=/tmp/comet-rust-python/bin/python" >> "$GITHUB_ENV"
echo "PYTHONPATH=$(/tmp/comet-rust-python/bin/python -c 'import sysconfig; print(sysconfig.get_path("purelib"))')" >> "$GITHUB_ENV"

- name: Check Cargo clippy
shell: bash
run: |
cd native
cargo clippy --color=never --all-targets --workspace -- -D warnings

- name: Check native Python UDF feature
shell: bash
run: |
cd native
cargo clippy --color=never -p datafusion-comet --all-targets --features python-udf -- -D warnings

- name: Check compilation
shell: bash
run: |
Expand Down Expand Up @@ -83,6 +99,13 @@ runs:
export LD_LIBRARY_PATH=${JAVA_HOME}/lib/server:${LD_LIBRARY_PATH}
RUST_BACKTRACE=1 cargo nextest run

- name: Run native Python UDF tests
shell: bash
run: |
cd native
export LD_LIBRARY_PATH=${JAVA_HOME}/lib/server:${LD_LIBRARY_PATH}
RUST_BACKTRACE=1 cargo nextest run -p datafusion-comet --lib --features python-udf python_udf

# The steps above lint and test the accounting allocator over the system allocator. Lint and
# test it over jemalloc too, since that is a different allocator backend and nothing else in CI
# builds it.
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -524,6 +524,7 @@ jobs:
org.apache.comet.exec.CometJoinSuite
org.apache.comet.exec.CometTypedDatasetSuite
org.apache.spark.sql.comet.CometMapInBatchSuite
org.apache.spark.sql.comet.CometArrowPythonUdfSuite
org.apache.spark.sql.execution.python.CometArrowPythonRunnerSuite
org.apache.comet.CometNativeSuite
org.apache.comet.CometConfSuite
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,7 @@ jobs:
org.apache.comet.exec.CometJoinSuite
org.apache.comet.exec.CometTypedDatasetSuite
org.apache.spark.sql.comet.CometMapInBatchSuite
org.apache.spark.sql.comet.CometArrowPythonUdfSuite
org.apache.spark.sql.execution.python.CometArrowPythonRunnerSuite
org.apache.comet.CometNativeSuite
org.apache.comet.CometConfSuite
Expand Down
36 changes: 30 additions & 6 deletions .github/workflows/pyarrow_udf_test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -44,12 +44,15 @@ jobs:
- name: Spark 4.0
maven_profiles: "-Pspark-4.0 -Pscala-2.13"
pyspark: "4.0.4"
native_python_udf: false
- name: Spark 4.1
maven_profiles: "-Pspark-4.1"
pyspark: "4.1.3"
native_python_udf: true
- name: Spark 4.2
maven_profiles: "-Pspark-4.2"
pyspark: "4.2.0"
native_python_udf: true
container:
# Pinned to the Debian 12 (bookworm) base so the system `python3` is 3.11. The default
# `amd64/rust` image is Debian 13 (trixie) which ships Python 3.13 and no python3.11 apt
Expand Down Expand Up @@ -77,19 +80,34 @@ jobs:
restore-keys: |
${{ runner.os }}-java-maven-

- name: Build Comet (debug, ${{ matrix.name }} / Scala 2.13)
run: |
cd native && cargo build
cd .. && ./mvnw -B install -DskipTests ${{ matrix.maven_profiles }}

- name: Install Python 3.11 and pip
run: |
apt-get update
apt-get install -y --no-install-recommends python3 python3-venv python3-pip
apt-get install -y --no-install-recommends python3 python3-dev python3-venv python3-pip
python3 -m venv /tmp/venv
/tmp/venv/bin/pip install --upgrade pip
/tmp/venv/bin/pip install "pyspark==${{ matrix.pyspark }}" "pyarrow>=14" pandas pytest

- name: Build Comet (debug, ${{ matrix.name }} / Scala 2.13)
env:
PYO3_PYTHON: /tmp/venv/bin/python
run: |
if [ "${{ matrix.native_python_udf }}" = "true" ]; then
(cd native && cargo build --features python-udf)
else
(cd native && cargo build)
fi
./mvnw -B install -DskipTests ${{ matrix.maven_profiles }}

- name: Run native Arrow UDF Scala suite
if: matrix.native_python_udf
env:
PYSPARK_PYTHON: /tmp/venv/bin/python
PYTHONPATH: /tmp/venv/lib/python3.11/site-packages
run: |
./mvnw -B test ${{ matrix.maven_profiles }} \
-Dsuites=org.apache.spark.sql.comet.CometArrowPythonUdfSuite

- name: Run PyArrow UDF pytest
env:
# Spark launches Python workers in a fresh subprocess and looks up `python3`
Expand All @@ -98,11 +116,17 @@ jobs:
# ModuleNotFoundError.
PYSPARK_PYTHON: /tmp/venv/bin/python
PYSPARK_DRIVER_PYTHON: /tmp/venv/bin/python
# The embedded interpreter does not inherit PySpark worker sys.path setup.
PYTHONPATH: /tmp/venv/lib/python3.11/site-packages
run: |
/tmp/venv/bin/python -m pytest -v \
spark/src/test/resources/pyspark/test_pyarrow_udf.py
/tmp/venv/bin/python -m pytest -v \
spark/src/test/resources/pyspark/test_pyarrow_udf_dictionary_shuffle.py
if [ "${{ matrix.native_python_udf }}" = "true" ]; then
/tmp/venv/bin/python -m pytest -v \
spark/src/test/resources/pyspark/test_native_arrow_udf_files.py
fi
/tmp/venv/bin/python -m pytest -v \
spark/src/test/resources/pyspark/test_pyarrow_udf_fuzz.py

Expand Down
8 changes: 7 additions & 1 deletion dev/ci/compute-changes.py
Original file line number Diff line number Diff line change
Expand Up @@ -136,11 +136,16 @@
],
# A real Python worker against each Spark 4.x Arrow runner. The list is
# deliberately narrow: the suite builds Comet three times, once per Spark
# version, and only the map-in-batch wiring can change its verdict.
# version, and covers map-in-batch wiring and the native scalar Arrow UDF.
"pyarrow_udf": [
"pom.xml",
"common/pom.xml",
"native/shuffle/src/spark_unsafe/row.rs",
"native/core/src/execution/python_udf.rs",
"native/core/src/execution/operators/arrow_python_udf.rs",
"spark/src/main/spark-4.1+/org/apache/spark/sql/comet/CometArrowEvalPythonExec.scala",
"spark/src/main/spark-4.1+/org/apache/spark/sql/comet/shims/ShimCometArrowEvalPythonExec.scala",
"spark/src/test/spark-4.1+/org/apache/spark/sql/comet/CometArrowPythonUdfSuite.scala",
"spark/pom.xml",
"spark/src/main/java/org/apache/comet/vector/**",
"spark/src/main/java/org/apache/spark/sql/comet/execution/shuffle/SpillWriter.java",
Expand All @@ -161,6 +166,7 @@
"spark/src/main/spark-4.x/org/apache/spark/sql/execution/python/CometArrowPythonRunnerBase.scala",
"spark/src/test/resources/pyspark/conftest.py",
"spark/src/test/resources/pyspark/test_pyarrow_udf.py",
"spark/src/test/resources/pyspark/test_native_arrow_udf_files.py",
"spark/src/test/resources/pyspark/test_pyarrow_udf_fuzz.py",
"spark/src/test/resources/pyspark/test_pyarrow_udf_dictionary_shuffle.py",
"spark/src/test/spark-3.5/org/apache/spark/sql/comet/CometMapInBatchSuite.scala",
Expand Down
13 changes: 8 additions & 5 deletions docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -439,11 +439,14 @@ diverge for several structural reasons:
allocation counters (`native_allocated`, `jemalloc_allocated`) see it. In a default build the C
dependencies are libzstd (`zstd-sys`, behind the Parquet `zstd` codec), libhdfs (`hdfs-sys`,
pulled in by the default `hdfs-opendal` feature), and the TLS stack used for cloud object stores
(`aws-lc-sys`). Building with the `jemalloc` or `mimalloc` feature adds the allocator itself
(`tikv-jemalloc-sys`, `libmimalloc-sys`). It is worth knowing which dependencies are _not_ C,
because several names suggest otherwise: the other Parquet codecs are pure Rust in this build,
`snap` for Snappy, `lz4_flex` for LZ4 and `zlib-rs` for gzip, as is `libbz2-rs-sys` despite its
name, so those allocations do pass through `GlobalAlloc` and are counted.
(`aws-lc-sys`). With the `python-udf` feature, the embedded Python interpreter and PyArrow also
allocate outside `GlobalAlloc`; their memory is absent from the executor's native `allocated`
figure, and Python working allocations are not reserved in Comet's pool. Building with the
`jemalloc` or `mimalloc` feature adds the allocator itself (`tikv-jemalloc-sys`,
`libmimalloc-sys`). It is worth knowing which dependencies are _not_ C, because several names
suggest otherwise: the other Parquet codecs are pure Rust in this build, `snap` for Snappy,
`lz4_flex` for LZ4 and `zlib-rs` for gzip, as is `libbz2-rs-sys` despite its name, so those
allocations do pass through `GlobalAlloc` and are counted.
- **Batches in flight across the FFI boundary.** Reservations stop at the operator that made them.
Imported JVM batches are reserved only while a reserving operator holds them, and exported native
batches have usually been released by the time the JVM receives them yet stay resident until the
Expand Down
9 changes: 5 additions & 4 deletions docs/source/user-guide/latest/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -141,10 +141,11 @@ natively on `BroadcastHashJoinExec` and `ShuffledHashJoinExec`. Existence sort-m

## Python and UDF

| Operator | Status | Notes |
| -------------------------------------------------- | ------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `MapInArrowExec`, `MapInPandasExec` | ⚠️ | Spark 4.0 and later. Experimental, disabled by default (`spark.comet.exec.pyarrowUDF.enabled`). See [PyArrow UDF Acceleration](pyarrow-udfs.md). |
| `ArrowEvalPythonExec`, `FlatMapGroupsInPandasExec` | 🔜 | Scalar `@pandas_udf` ([#5386](https://github.com/apache/datafusion-comet/issues/5386)) and grouped `applyInPandas` ([#5123](https://github.com/apache/datafusion-comet/issues/5123)) fall back to Spark. |
| Operator | Status | Notes |
| -------------------------------------------------------------- | ------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `ArrowEvalPythonExec` for scalar `@arrow_udf` (Spark 4.1+) | ⚠️ | Experimental native PyO3 path, requires the `python-udf` feature and `spark.comet.exec.nativeArrowPythonUDF.enabled`. See [PyArrow UDF Acceleration](pyarrow-udfs.md). |
| `MapInArrowExec`, `MapInPandasExec` | ⚠️ | Spark 4.0 and later. Experimental, disabled by default (`spark.comet.exec.pyarrowUDF.enabled`). See [PyArrow UDF Acceleration](pyarrow-udfs.md). |
| Other `ArrowEvalPythonExec` types, `FlatMapGroupsInPandasExec` | 🔜 | Scalar `@pandas_udf` ([#5386](https://github.com/apache/datafusion-comet/issues/5386)) and grouped `applyInPandas` ([#5123](https://github.com/apache/datafusion-comet/issues/5123)) fall back to Spark. |

## See also

Expand Down
97 changes: 89 additions & 8 deletions docs/source/user-guide/latest/pyarrow-udfs.md
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,84 @@ spark.comet.exec.pyarrowUDF.enabled=true

The default is `false` while the feature stabilizes.

### Native scalar Arrow UDFs (Spark 4.1+)

Scalar `@arrow_udf` in Spark 4.1 and later can run inside Comet's Rust execution pipeline when
the native library is built with the `python-udf` Cargo feature and this separate option is enabled:

```
spark.comet.exec.nativeArrowPythonUDF.enabled=true
```

Build the native library with a Python interpreter that matches the major and minor version used
by the executors' PySpark workers. That interpreter needs its development headers (for example,
`python3-dev` on Debian and Ubuntu) and a shared `libpython` (`libpython3.x.so` on Linux). A Python
built with pyenv may need `PYTHON_CONFIGURE_OPTS="--enable-shared"` when it is installed. For example:

```sh
PYO3_PYTHON=/path/to/python make release COMET_FEATURES=python-udf
```

Install the matching shared `libpython` on every executor and make it discoverable by the dynamic
linker. On Linux, the JVM loads `libcomet` with local symbols, so Comet promotes the already-loaded
shared `libpython` to the global namespace before importing Python extensions such as PyArrow. A
statically linked Python cannot provide those symbols this way, and importing PyArrow can fail with
an undefined-symbol error. A library built with `python-udf` depends on `libpython` as soon as
`libcomet` is loaded: if an executor cannot find it, **Comet itself fails to load**, even when a
query does not use an Arrow UDF.

Comet passes each argument as a `pyarrow.Array` through the Arrow C Data Interface, invokes the
pickled Python function with PyO3, and appends the result array to the input batch. It checks the
result length and safely casts it to the declared return type, matching Spark's scalar Arrow UDF
serializer. It splits larger input batches according to
`spark.sql.execution.arrow.maxRecordsPerBatch` and
`spark.sql.execution.arrow.maxBytesPerBatch`. A native worker is created per partition.

Each partition unpickles its own callable. Imported modules and their global state are shared by
concurrent tasks in the executor's embedded Python interpreter.

The executor's embedded Python must be able to import `pyspark`, `pyarrow`, and the user's Python
modules. Spark serializes the callable and a PySpark return type with `pyspark.cloudpickle`, so
`pyspark` is required even when the callable itself only uses PyArrow. Build the `python-udf`
feature against the same Python major/minor version used by PySpark workers. Install those packages
into that Python environment and ensure the executor process can find them through the embedded
interpreter's `sys.path` (for example, by setting `PYTHONPATH` before launching the executor).
`PYSPARK_PYTHON` selects the external worker executable; it does not select or configure the
embedded interpreter. The worker-only `pyspark.zip` path is not automatically added to it. The
feature and config are disabled by default. Without either, `ArrowEvalPythonExec` stays on Spark's
normal path.

The initial native path accepts scalar `@arrow_udf` calls with regular or named arguments and
multiple independent UDFs in one `ArrowEvalPythonExec`. Chained Python UDFs, broadcast variables,
Python includes, per-function environment overrides other than Spark's default
`PYTHONHASHSEED=0`, and `spark.sql.execution.arrow.useLargeVarTypes=true` stay on Spark's path.
UDFs that capture a PySpark accumulator also stay on Spark's worker path so their task updates
reach the driver.
Queries also stay on Spark's Python worker path when the Spark context has files added through
`addPyFile` or `addFile`, because the embedded interpreter does not receive Spark's per-task file
setup.
The embedded interpreter starts with the same default hash seed as Spark's Python workers.
Iterator Arrow UDFs,
ordinary `udf(..., useArrow=True)`, scalar pandas UDFs, and `mapInArrow` are separate execution
types; `mapInArrow` retains the columnar runner described above.

The native path accepts boolean, byte, short, integer, long, float, double, plain string, binary,
decimal, date, and timestamp without time zone. Other input or result types and an
enabled `spark.sql.pyspark.udf.profiler` stay on Spark's path. Spark labels `TimestampType` with the
session time zone, while Comet uses UTC; nested Arrow field names can also differ. This allow-list
keeps types with unverified Arrow schemas on Spark's path. `TimeType` also stays on Spark's path:
Spark 4.2's Arrow UDF row converter rejects it even though PySpark can describe its Arrow type.

The embedded interpreter is shared by tasks. Pure Python code contends on its global interpreter
lock, so multiple partitions may be slower than Spark's separate Python workers; PyArrow kernels
that release the lock can still run concurrently. Python execution stays synchronous on JVM input
paths and hands off other async tasks when it runs on a Tokio worker. `pyspark.TaskContext.get()`
returns `None` inside a native UDF. A native extension crash or `os._exit` terminates the executor
process. The embedded Python interpreter and PyArrow allocate outside Comet's memory pool. Those
allocations are also absent from the executor's `Comet native memory usage: allocated` figure and
are not limited by `spark.executor.pyspark.memory`. Budget them in executor memory overhead in
addition to the [memory log estimate](tuning/memory.md#sizing-the-overhead-from-the-memory-usage-log).

### Relationship to Spark's PySpark Arrow conversion conf

`spark.comet.exec.pyarrowUDF.enabled` is **not** the same as PySpark's
Expand All @@ -91,12 +169,14 @@ worker. Both confs can be set independently.

## Supported APIs

| PySpark API | Spark Plan Node | Supported |
| -------------------------------- | --------------------------- | --------- |
| `df.mapInArrow(func, schema)` | `MapInArrowExec` | Yes |
| `df.mapInPandas(func, schema)` | `MapInPandasExec` | Yes |
| `@pandas_udf` (scalar) | `ArrowEvalPythonExec` | Not yet |
| `df.applyInPandas(func, schema)` | `FlatMapGroupsInPandasExec` | Not yet |
| PySpark API | Spark Plan Node | Supported |
| -------------------------------- | --------------------------- | ------------------------ |
| `df.mapInArrow(func, schema)` | `MapInArrowExec` | Yes |
| `df.mapInPandas(func, schema)` | `MapInPandasExec` | Yes |
| scalar `@arrow_udf` (Spark 4.1+) | `ArrowEvalPythonExec` | Experimental native path |
| `udf(..., useArrow=True)` | `ArrowEvalPythonExec` | Not yet |
| `@pandas_udf` (scalar) | `ArrowEvalPythonExec` | Not yet |
| `df.applyInPandas(func, schema)` | `FlatMapGroupsInPandasExec` | Not yet |

## Example

Expand Down Expand Up @@ -177,8 +257,9 @@ on the unoptimized path.

## Limitations

- The optimization currently applies only to `mapInArrow` and `mapInPandas`. Scalar pandas UDFs
(`@pandas_udf`) and grouped operations (`applyInPandas`) are not yet supported.
- The columnar Python runner applies to `mapInArrow` and `mapInPandas`. The separate native path
applies to scalar `@arrow_udf` on Spark 4.1+. Scalar pandas UDFs (`@pandas_udf`) and grouped
operations (`applyInPandas`) are not yet supported.
- The optimization requires Arrow data on the input side. If a shuffle sits between the upstream
Comet operator and the Python UDF, use Comet's columnar shuffle for the optimization to apply.
Both the `jvm` and `native` shuffle modes can feed `CometMapInBatch`. Set
Expand Down
Loading
Loading