Skip to content

feat: execute scalar Arrow Python UDFs in native Comet - #6130

Open
viirya wants to merge 13 commits into
apache:mainfrom
viirya:codex/native-arrow-python-udf
Open

viirya wants to merge 13 commits into
apache:mainfrom
viirya:codex/native-arrow-python-udf

Conversation

@viirya

@viirya viirya commented Sep 23, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6129.

Rationale for this change

Starting with Spark 4.1, scalar @arrow_udf expressions fall back at ArrowEvalPythonExec, interrupting Comet's native pipeline. An opt-in in-process path keeps Arrow batches in Comet and avoids the Spark Python worker IPC path for supported UDFs.

What changes are included in this PR?

  • Add an optional python-udf Cargo feature and a native ArrowPythonUdfExec that invokes Python through PyO3 and exchanges arrays through the Arrow C Data interface. Python evaluation uses block_in_place, preserving synchronous polling for JVM scans while allowing Tokio to hand off other work on native-source paths.
  • Serialize verified Spark 4.1+ scalar Arrow UDF types into the native plan. An allow-list covers boolean, byte, short, integer, long, float, double, plain string, binary, decimal, date, and timestamp without time zone. Other types, including timezone-aware timestamps, complex types, variant, interval, and time, fall back. Spark 4.2's Arrow UDF row converter itself rejects TimeType, so the native path preserves that behavior. Profiler-enabled UDFs also fall back. Earlier Spark versions retain their current behavior.
  • Add spark.comet.exec.nativeArrowPythonUDF.enabled, defaulting to false, plus documentation of embedded Python setup and process-level limitations.
  • Build the feature and run the native Scala suite in the Spark 4.1 and 4.2 PyArrow matrix. Feature-gated Rust clippy and unit tests run in the common Rust test action with embedded Python dependencies and the JVM library path, so generic native changes also exercise them.
  • Move the benchmark into the benchmark package and add PyArrow-kernel and Python-loop modes.

The callable still executes in Python. This PR supports scalar Arrow UDFs only; regular Python UDFs and pandas UDFs remain outside this path. All tasks share an embedded interpreter, so Python code that holds the GIL can be slower than Spark's worker-per-task execution.

How are these changes tested?

  • cargo clippy -p datafusion-comet --all-targets --features python-udf -- -D warnings passed; five Rust unit tests passed with the feature enabled.
  • Spark 4.1 and 4.2 native integration suites passed with a feature-enabled library. The test command uses pyspark.cloudpickle and a real PySpark return type. End-to-end cases compare Spark and native results for every allowed scalar type with a null, alongside multiple functions and fallback cases for unsupported schemas, profiler, and configuration settings.
  • Prettier, Spotless, cargo fmt --check, and CI configuration checks passed. Spark 3.5 and 4.0 compatibility profiles compiled earlier in this PR.
  • After replacing spawn_blocking with block_in_place, a Spark 4.1.3 benchmark of 50 million rows with pyarrow.compute.negate and two partitions gave medians of 1.425 s for Spark and 0.398 s for native (3.58x faster), from three measured iterations after warmup. The query verifies its aggregate checksum.
  • The same revision's Python-loop benchmark used 3.2 million rows, 32 partitions, and local[16]: Spark median 0.143 s and native median 0.492 s (native 3.44x slower). The prior spawn_blocking version took 1.055 s on this native case, consistent with busy polling on the JVM scan path. The feature remains opt-in because GIL-bound work can still favor Spark workers.

The run-spark-4.1-tests and run-pyarrow-udf-tests labels request broader CI coverage before merge.

@github-actions github-actions Bot added the enhancement New feature or request label Sep 23, 2026
@viirya viirya added run-pyarrow-udf-tests Run the PyArrow UDF tests on this pull request instead of waiting for the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue labels Sep 23, 2026
@viirya
viirya marked this pull request as draft September 23, 2026 03:04
@viirya
viirya marked this pull request as ready for review September 23, 2026 07:48
case udf if udf.func.pythonIncludes != null && !udf.func.pythonIncludes.isEmpty =>
"Arrow UDF Python includes are not supported in-process"
case udf if udf.func.envVars != null && !udf.func.envVars.isEmpty =>
"Arrow UDF Python environment overrides are not supported in-process"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This rejects ordinary @arrow_udf calls because PySpark always adds PYTHONHASHSEED to envVars. Please handle Spark's default environment and add a PySpark test asserting native execution. The current tests use an empty map and miss this case.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for catching this. Fixed in 4830cee65: native execution now accepts Spark’s default PYTHONHASHSEED=0, while custom values still fall back to Spark. The embedded interpreter uses hash seed 0 as well. I added a test using a real PySpark @arrow_udf that asserts native execution is selected and compares its results, including hash(), with Spark.

Comment thread native/core/src/execution/python_udf.rs Outdated
"_import_from_c",
(
&raw const ffi_array as Py_uintptr_t,
&raw const ffi_schema as Py_uintptr_t,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please use let mut and &raw mut for these structs, including ffi_return_type below. PyArrow's import writes their release fields, which is undefined behavior through these &raw const pointers.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 4830cee65. The input FFI_ArrowArray and FFI_ArrowSchema, as well as ffi_return_type, are now mutable and passed with &raw mut so PyArrow can update their release fields. The feature-enabled Rust tests pass.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Scalar Arrow UDFs fall back to Spark’s Python worker, interrupting native execution.
  • Design approach: Add an opt-in PyO3 operator using Arrow C Data transfers to avoid worker IPC.
  • Correctness / compatibility analysis: Compared Spark 4.1.3 and 4.2.0 sources. Found one new P2 issue: native execution ignores Spark’s configured UDF input batch limit. The existing environment-map and immutable FFI-pointer concerns remain substantiated and unresolved.
  • Key design decisions: Feature gating, version shims, and a conservative type allow-list keep the implementation contained. Shared interpreter state and GIL contention are documented limitations.
  • Implementation sketch: Scala serializes functions and arguments. Rust creates per-partition callables, evaluates batches, validates results, and appends output columns.
  • Behavioral changes worth calling out: Disabled by default. Python allocations remain outside Comet’s memory pool. Reported benchmarks show different outcomes for PyArrow kernels and GIL-bound Python work.
  • Suggested improvements: Preserve Spark’s configured Arrow batch limit before invoking Python.

Reviewed all 28 changed files at b8cc409b13d52a0b8744949ece8166d215ff5511 against 971380971064082215744f4fc4d86177bdc416c7. PR confirmed non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 34 successful checks, 13 skipped, none failed or pending. Linux builds, Rust tests, Spark 4.1 SQL suites, and PyArrow jobs for Spark 4.0–4.2 passed. Logs confirm three native Scala cases passed on each of Spark 4.1 and 4.2, plus all five Python bridge Rust tests.

Validation: All five bridge tests passed locally in an isolated harness using the exact source. A bounded operator reproduction and an actual Spark 4.1.3 query confirmed the batch-limit mismatch. No full local Comet JNI/Scala build, macOS validation, or benchmark rerun was performed. The two existing threads are excluded from new findings.

// Keep the JVM scan path synchronous so its Pending loop does not spin while
// Python runs. On a tokio worker, this hands its other tasks to another worker.
tokio::task::block_in_place(|| {
Self::evaluate_batch(&specs, &workers, Arc::clone(&schema), batch?)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could this path preserve spark.sql.execution.arrow.maxRecordsPerBatch before invoking Python? With native execution selected, a cap of 2, and a four-row child batch, Spark calls the UDF on batches of at most two rows, but this code forwards all four rows. A UDF that rejects batches larger than two succeeds in Spark and fails here with ValueError: configured batch cap exceeded: 4 > 2. This breaks workloads that use Spark’s cap to satisfy Python library or memory limits. Serialize the configured limit and slice all input columns at matching boundaries before evaluation, or fall back when the limit cannot be honored.

Evidence: Ran Spark 4.1.3 with spark.sql.execution.arrow.maxRecordsPerBatch=2 and spark.range(1, 5, 1, 1).select(arrow_udf(capped_identity, LongType())('id')), where capped_identity(a) raises ValueError when len(a) > 2 and otherwise returns a. Spark returned [1, 2, 3, 4]. An isolated Rust harness importing the exact-head ArrowPythonUdfExec and bridge source used the same PySpark-cloudpickled callable and return type. One four-row child batch produced the stated error; two two-row batches returned all four rows successfully. Spark’s BatchedPythonArrowInput.writeSizedBatch enforces the cap, while the new protobuf and operator carry no such limit. Reproduction: /tmp/comet-6130-spark-reference.py and /tmp/comet-6130-bridge-test/tests/batch_limit.rs.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 4830cee65. We now pass arrowMaxRecordsPerBatch to the Rust operator and slice each input batch at that limit before invoking Python. I added a real PySpark test with the limit set to 2; it asserts that native execution is selected and that Python sees the same batch sizes as Spark.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Scalar Arrow UDFs use Spark’s Python workers, interrupting Comet’s native pipeline.
  • Design approach: Add an opt-in PyO3 operator that exchanges arrays through Arrow C Data.
  • Correctness / compatibility analysis: Compared Spark 4.1.3 and 4.2.0 sources. Found two introduced P2 issues: the byte-based input batch limit is ignored, and the new regression test uses an unsuitable Spark 4.2 reference path. The three previously reported concerns are addressed at this head.
  • Key design decisions: Feature gating, version shims, and the scalar type allow-list keep the implementation contained. The operator and bridge have clear responsibilities without unnecessary abstraction.
  • Implementation sketch: Scala serializes callables, arguments, return types, and the row cap. Rust creates per-partition callables, slices batches, evaluates Python, validates results, and appends columns.
  • Behavioral changes worth calling out: Disabled by default. Shared interpreter state, GIL contention, and Python allocations outside Comet’s memory pool are documented. Reported benchmarks distinguish PyArrow kernels from slower GIL-bound work. I did not rerun benchmarks.
  • Suggested improvements: Preserve Spark’s byte-based batch limit and isolate the regression test’s Spark reference from Comet columnar input.

Reviewed all 29 changed files at 4830cee650fb15a4540b3cd4c2b6aa322ddd8921 against 971380971064082215744f4fc4d86177bdc416c7. Confirmed non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 32 successful checks, 13 skipped, and 2 failed. The Spark 4.2 PyArrow job fails the new batch-limit test, causing Required Checks to fail. Linux builds, Rust tests, Spark 4.1 SQL suites, and Spark 4.0/4.1 PyArrow jobs passed. Logs confirm five bridge Rust tests and three native Scala cases on each of Spark 4.1 and 4.2 passed.

Validation: Five bridge tests passed locally in an isolated harness importing the unmodified exact-head sources. A focused operator test verified the row-cap fix and reproduced the byte-cap failure against an actual Spark 4.1.3 query. No full local Comet JNI/Scala build, local Spark 4.2 execution, or macOS validation was performed.

val native = OperatorOuterClass.ArrowPythonUdf
.newBuilder()
.addAllFunctions(functions.map(_.get).asJava)
.setMaxRecordsPerBatch(op.conf.arrowMaxRecordsPerBatch)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could this preserve spark.sql.execution.arrow.maxBytesPerBatch alongside the row cap? With a byte cap of 16, a row cap of 10000, and four long values, Spark sends two-row batches to Python. Native execution sends all four because only the row cap reaches the operator. A UDF that rejects a.nbytes > 16 therefore succeeds in Spark but fails natively with ValueError: configured byte cap exceeded: 32 > 16. This breaks workloads using the byte limit to satisfy Python library or memory constraints. Serialize and enforce the byte limit on UDF arguments, or fall back when it cannot be honored.

Evidence: Ran Spark 4.1.3 on spark.range(1, 5, 1, 1) with arrow.maxRecordsPerBatch=10000 and arrow.maxBytesPerBatch=16. The UDF checks a.nbytes > 16 and otherwise returns a. Spark returned [1, 2, 3, 4], with observed batch lengths [2, 2, 2, 2]. An isolated Rust harness importing this head’s operator and bridge used the same PySpark-cloudpickled command and one four-row child batch. It produced the stated 32-byte error. Setting the native row cap to 2 succeeded. Spark’s BatchedPythonArrowInput.writeSizedBatch enforces both limits. Reproduction files: /tmp/comet-6130-4830-harness/reference.py and /tmp/comet-6130-4830-harness/tests/operator.rs.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the detailed reproduction. Fixed in c73d8b599: Scala now passes arrowMaxBytesPerBatch to the native operator. Rust splits all UDF arguments at the same boundary when either the row or byte limit is reached, counting Arrow values, offsets, and validity buffers. I added Rust tests and a real PySpark @arrow_udf test with a 16-byte limit and four long values. The Spark 4.1 and 4.2 PyArrow CI jobs pass.

spark.conf.set("spark.comet.sparkToColumnar.enabled", "true")
try:
spark.conf.set("spark.comet.exec.nativeArrowPythonUDF.enabled", "false")
spark_rows = source.select(batch_length("id")).collect()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could the reference query run with Comet disabled, restoring it before the native query? Disabling only nativeArrowPythonUDF.enabled leaves sparkToColumnar.enabled=true. On Spark 4.2, the reference consequently uses Spark’s columnar Python runner, which forwards the four-row input batch without applying the row cap. The reference returns [4, 4, 4, 4], while the corrected native operator returns [2, 2, 2, 2]. This makes the new test fail and blocks Required Checks. Use a row-based Spark reference while retaining the native-plan and two-row-cap assertions.

Evidence: Exact-head job https://github.com/apache/datafusion-comet/actions/runs/36100068127/job/107961195675 fails at line 142 with Row(batch_length(id)=2) != Row(batch_length(id)=4). It reports 1 failed and 133 passed. Spark v4.2.0’s ArrowEvalPythonExec selects columnar execution for a columnar child, and ColumnarArrowPythonInput.writeRowByRow writes the whole batch without the BatchedPythonArrowInput limits. The corresponding Spark 4.1 job passes this test.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, I confirmed the Spark 4.2 failure. The reference query now runs with spark.comet.enabled=false and checks that its plan contains no Comet operator. The test restores Comet before the native query and retains the native-plan and two-row batch assertions. The Spark 4.2 PyArrow CI job now passes.

@viirya

viirya commented Sep 25, 2026

Copy link
Copy Markdown
Member Author

Addressed both review comments in c73d8b599: native scalar Arrow UDFs now honor the Arrow byte batch limit, and the Spark reference test runs without Comet. PyArrow UDF CI passes on Spark 4.0, 4.1, and 4.2.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Scalar Arrow UDFs use Spark’s Python workers, interrupting Comet’s native pipeline.
  • Design approach: Add an opt-in PyO3 operator that exchanges arrays through Arrow C Data.
  • Correctness / compatibility analysis: Compared Spark 4.1.3 and 4.2.0 sources. Found one new P2 issue: empty native batches invoke Python when Spark skips evaluation. The five previously reported concerns are addressed at this head.
  • Key design decisions: Feature gating, version shims, and a conservative type allow-list contain the scope. The operator and Python bridge have clear responsibilities.
  • Implementation sketch: Scala serializes functions, arguments, return types, and batch limits. Rust creates per-partition callables, slices inputs, validates Python results, and appends output columns.
  • Behavioral changes worth calling out: Disabled by default. Shared interpreter state, GIL contention, and Python allocations outside Comet’s memory pool are documented. Reported benchmarks distinguish PyArrow-kernel gains from GIL-bound slowdowns.
  • Suggested improvements: Skip empty batches before argument/UDF evaluation and add coverage for sampling that produces an empty batch.

Reviewed all 29 changed files at c73d8b599e82ae81a484f3d3cf5287b4f9a07217 against 971380971064082215744f4fc4d86177bdc416c7. Confirmed non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 23 successful checks, 11 skipped, 10 in progress, none failed. PyArrow jobs passed on Spark 4.0–4.2. Logs confirm eight native feature tests and three native Scala cases on each of Spark 4.1 and 4.2 passed. Broader Linux and Spark 4.1 SQL checks remain pending.

Validation: Eight feature unit tests passed locally in an isolated harness importing the unmodified head sources. Three focused operator tests verified both batch caps and reproduced the empty-batch failure, including SampleExec, against an actual Spark 4.1.3 reference query. No full local Comet JNI/Scala build, local Spark 4.2 execution, macOS validation, or benchmark rerun was performed.

}
let batch = batch.as_ref()?;
let args = args.as_ref()?;
if offset == batch.num_rows() && (offset != 0 || emitted_empty_batch) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could zero-row input batches be skipped before evaluating arguments or calling Python? On the first poll, this condition deliberately allows an empty batch through. With native Arrow UDF execution and Spark-to-columnar conversion enabled, spark.range(1, 5, 1, 1).sample(False, 0.01, 42).select(f("id")) reaches this case because native SampleExec emits an empty batch. If f rejects empty arrays, Spark returns [] without invoking it, but native execution fails. This breaks valid UDFs using libraries that require nonempty input. Skip empty batches and consider adding this sampling case as a regression test.

Evidence: Ran Spark 4.1.3 with f = arrow_udf(nonempty, LongType()), where nonempty(a) raises ValueError("unexpected empty UDF input") when len(a) == 0 and otherwise returns a. The sampling query returned []. An isolated Rust harness importing this head’s unmodified SampleExec, ArrowPythonUdfExec, and Python bridge, with unchanged sampler/RNG definitions, used the same PySpark-cloudpickled callable. Sampling [1, 2, 3, 4] with bounds [0, 0.01) and seed 42 produced one zero-row batch, then failed with that ValueError. Spark’s BatchedPythonArrowInput writes only when nextBatchStart.hasNext. Reproduction files: /tmp/comet-6130-c73d-review/reference.py and /tmp/comet-6130-c73d-review/tests/operator.rs.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the precise reproduction. Fixed in 67a9f776a: the native operator now skips zero-row input batches before evaluating arguments or calling Python. I added a regression test using the same sampling case and a real @arrow_udf that rejects empty arrays. It verifies that Spark and native execution both return [], with the native operator selected. The PyArrow UDF jobs pass on Spark 4.1 and 4.2.

@andygrove

Copy link
Copy Markdown
Member

This is a light fully automated review since there are so many PRs open.

Unlike its sibling operators, CometArrowEvalPythonExec (spark/src/main/spark-4.1+/org/apache/spark/sql/comet/CometArrowEvalPythonExec.scala:152) doesn't override stringArgs, equals or hashCode. Spark's default stringArgs is productIterator, so argString prints nativeOp through protobuf's toString, which dumps the whole native subtree. For spark.read.parquet("s3a://...").select(f("x")) that includes the pickled command bytes and the child scan's object_store_options, which NativeConfig.extractObjectStoreOptions fills with every fs.s3a.* key, fs.s3a.secret.key included. That text lands in explain(), the SQL UI and the event log, although CometNativeScanExec itself prints only output. The default equals and hashCode also compare nativeOp, whose plan_id differs between two otherwise identical subtrees, so an exchange above a native UDF can't be reused and a self-join runs the Python function once per side. Could this override the three methods the way CometProjectExec does, keeping the UDFs in the identity so two functions with the same result type still differ? A test asserting the plan string stays free of the command would guard against a regression.

docs/source/user-guide/latest/pyarrow-udfs.md:153 says PyArrow allocations are outside Comet's memory pool. They are also missing from the allocated figure in the executor's memory usage log, which counts only Rust's global allocator, even though the UDF output columns then flow through native operators that may reserve them. Sizing spark.executor.memoryOverhead from allocated - reserved, as the tuning guide describes, would come up short by whatever the interpreter and PyArrow hold. Could the guide say to add that on top, and could the non-Rust allocations list at docs/source/contributor-guide/memory_management.md:364 mention the embedded interpreter and PyArrow's allocator?

@viirya

viirya commented Sep 28, 2026

Copy link
Copy Markdown
Member Author

Thanks for catching both issues. In 67a9f776a, CometArrowEvalPythonExec.stringArgs no longer renders nativeOp, so the pickled command and native subtree are absent from the plan string. equals and hashCode compare the Arrow UDF payload while ignoring the outer plan_id. A test checks that the command is absent from the plan, changing only the plan ID preserves equality, and changing the command does not.

I also documented that embedded Python and PyArrow allocations are excluded from Rust's allocated figure and need additional executor memory overhead. After merging upstream in de71a8e78, that guidance lives in the new tuning/memory.md page and covers the JVM Arrow metric's treatment of PyArrow buffers imported from native. Preflight and the PyArrow UDF jobs on Spark 4.0, 4.1, and 4.2 pass; the Spark 4.1 SQL gate is still running.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Scalar Arrow UDFs use Spark’s Python workers, interrupting Comet’s native pipeline.
  • Design approach: Add an opt-in PyO3 operator that exchanges arrays through Arrow C Data.
  • Correctness / compatibility analysis: Compared Spark 4.1.3 and 4.2.0 sources. Found one new P2 issue: scalar-subquery arguments disappear from Spark’s expression traversal and Comet’s subquery registration. Earlier environment, FFI-pointer, batching, and empty-input concerns are addressed.
  • Key design decisions: Feature gating, version shims, and the conservative scalar-type allow-list contain the scope. The operator and Python bridge have clear responsibilities without unnecessary abstraction.
  • Implementation sketch: Scala serializes callables, arguments, types, and batch limits. Rust creates per-partition callables, slices inputs, validates Python results, and appends output columns.
  • Behavioral changes worth calling out: Disabled by default. Shared interpreter state, GIL contention, and Python allocations outside Comet’s memory accounting are documented. Reported benchmarks distinguish PyArrow-kernel gains from GIL-bound slowdowns. I did not rerun benchmarks.
  • Suggested improvements: Retain Catalyst UDF arguments for subquery preparation and registration, and finish the existing exchange-identity fix.

The existing exchange-reuse concern in #6130 (comment) remains substantiated. Equality now excludes the outer plan_id, but still compares generated argument expr_id values inside nativeOp.getArrowPythonUdf. A focused Spark reuse-rule reproduction produced one reused exchange for Spark, zero for otherwise identical native nodes with different generated IDs, and one after normalizing those IDs. This duplicates existing feedback and is excluded from new findings.

Reviewed all 31 changed files at de71a8e78d5f3570339d3fd99f84cbebcedd7738 against 65b334bfdd25196091a42d773bc3e61d19212e52. Confirmed non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 34 successful checks, 13 skipped, none failed or pending. Linux builds, Rust tests, Spark 4.1 SQL suites, and Spark 4.0–4.2 PyArrow jobs passed. Logs confirm eight native feature tests and four native Scala cases on each of Spark 4.1 and 4.2 passed.

Validation: Eight feature unit tests and two focused operator tests passed locally using the unmodified head sources. An actual Spark 4.1.3 reference query and a rebuilt isolated Scala harness reproduced the subquery-registration failure and exchange-reuse issue. The Scala harness preserves the reviewed case class and relevant Comet methods but stubs unrelated execution scaffolding. No full local Comet JNI/Scala build, local Spark 4.2 execution, macOS validation, or benchmark rerun was performed.

override val nativeOp: Operator,
override val originalPlan: SparkPlan,
override val output: Seq[Attribute],
resultAttrs: Seq[Attribute],

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could this node retain op.udfs or their Catalyst arguments as case-class fields? With native Arrow UDF execution and Spark-to-columnar conversion enabled, SELECT arrow_negate(id + (SELECT max(id) FROM range(8))) FROM range(4), where arrow_negate wraps pyarrow.compute.negate, should return [-7, -8, -9, -10]. The arguments are serialized successfully, but this node exposes only output attributes to QueryPlan.expressions. Spark does not traverse originalPlan for expressions, so Comet’s collectSubqueries finds nothing and native evaluation reaches CometScalarSubquery without a registered value, failing with Subquery ... not found for plan .... This breaks valid scalar-subquery arguments when the feature is enabled. Preserve the argument expressions for preparation and registration, or fall back for this case.

Evidence: Spark 4.1.3 executed the query and returned [-7, -8, -9, -10], with one ScalarSubquery inside ArrowEvalPythonExec. A rebuilt Scala harness using the exact-head CometArrowEvalPythonExec case class, unchanged Comet subquery collection/registration methods, and the current CometScalarSubquery.java reported: Spark has 1, native collector has 0; lookup=Subquery 0 not found for plan 6130. Production buildNativeContext uses that collector, and CometExecRDD registers only its results. Reproduction: /tmp/comet-6130-final-validation/reference.py and reference.log. This is an isolated lifecycle reproduction, not a full native query run.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed. CometArrowEvalPythonExec now retains op.udfs as a case-class field, so Spark exposes the UDF arguments through QueryPlan.expressions and Comet can collect and register the scalar subquery. I added a native execution regression using your id + (SELECT max(id) FROM range(8)) example and checking [-7, -8, -9, -10] in b0dab2f53.

@rich7420 rich7420 left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks for the update!

Comment on lines +85 to +86
case udf if udf.func.pythonIncludes != null && !udf.func.pythonIncludes.isEmpty =>
"Arrow UDF Python includes are not supported in-process"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

addPyFile with a plain .py file and addFile both leave pythonIncludes empty, so they pass this check even though the embedded interpreter hasn't set up Spark's files.

I tested both cases with Spark 4.1.3 and this head's native bridge. The UDFs succeed in Spark, but the added module fails to load in the bridge with ModuleNotFoundError, and reading the data file through SparkFiles.get() raises AssertionError. The native checks used a component harness. I haven't run the full Comet queries.

Could we keep these cases on Spark's worker path until file setup is supported, and add regression tests for both?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed that pythonIncludes does not cover a plain .py passed to addPyFile or a file passed to addFile. Native planning now checks SparkContext.listFiles() and keeps these queries on Spark's Python worker path. I added separate regression cases for addPyFile(.py) and addFile with SparkFiles.get(), each using a fresh Spark context, and documented the fallback in b0dab2f53.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Scalar Arrow UDFs use Spark’s Python workers, interrupting Comet’s native pipeline.
  • Design approach: Add an opt-in PyO3 operator that exchanges arrays through Arrow C Data.
  • Correctness / compatibility analysis: Compared Spark 4.1.3 and 4.2.0 sources. Found one introduced P2 issue: embedded execution silently loses PySpark accumulator updates. Previously reported environment, FFI-pointer, batching, empty-input, subquery, file-fallback, and exchange-identity concerns are addressed at this head.
  • Key design decisions: Feature gating, version shims, and a conservative type allow-list contain the scope. The operator and Python bridge have clear responsibilities without unnecessary abstraction.
  • Implementation sketch: Scala serializes callables, arguments, return types, and batch limits. Rust creates per-partition callables, slices inputs, validates Python results, and appends output columns.
  • Behavioral changes worth calling out: Disabled by default. Shared interpreter state, GIL contention, and Python allocations outside Comet’s memory accounting are documented. Reported benchmarks distinguish PyArrow-kernel gains from GIL-bound slowdowns. I did not rerun benchmarks.
  • Suggested improvements: Preserve task-scoped accumulator updates and cleanup, or retain Spark execution for accumulator-bearing UDFs until supported.

Reviewed all 32 changed files at b0dab2f5373e2a39451785f53efe61e9ec885bc4 against 65b334bfdd25196091a42d773bc3e61d19212e52. Confirmed non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 34 successful checks, 14 skipped, none failed or pending. Linux builds, Rust tests, Spark 4.1 SQL suites, and Spark 4.0–4.2 PyArrow jobs passed. Logs confirm eight native feature tests and four native Scala cases on each of Spark 4.1 and 4.2 passed, including the latest subquery and identity assertions. Both added-file fallback cases passed on both versions.

Validation: Eight feature unit tests and two focused operator tests passed locally in an isolated harness importing the unmodified head sources. An actual Spark 4.1.3 reference query and the same real PySpark-serialized callable reproduced the accumulator discrepancy in that harness. No full local Comet JNI/JVM build, local Spark 4.2 execution, macOS validation, or benchmark rerun was performed.

}
}
let pickle = py.import("pickle").map_err(python_error)?;
let loaded = pickle

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could accumulator-bearing UDFs remain on Spark’s worker path until embedded execution can forward their updates? With acc = sc.accumulator(0) and a scalar @arrow_udf('long') that calls acc.add(len(values)) before returning values, collecting four rows should update acc.value to 4. This path successfully unpickles and executes the callable, but only updates the embedded interpreter’s _accumulatorRegistry. It never performs Spark’s accumulator reporting, so the driver value remains 0 despite successful evaluation. Subsequent operator instances also reuse the retained accumulator state. This silently loses application counters. Forward updates through Spark’s accumulator mechanism with task isolation and cleanup, or fall back for these UDFs.

Evidence: Spark 4.1.3 executed spark.range(1, 5, 1, 1).select(f('id')).collect() and returned [1, 2, 3, 4] with accumulator value 4. A rebuilt harness importing this exact head’s unmodified ArrowPythonUdfExec and ArrowPythonUdf used a real PySpark command with eligible metadata: PYTHONHASHSEED=0, no includes, and no broadcasts. Two operator instances returned the correct rows while the embedded registry retained values 4 then 8. The driver accumulator stayed 0. Spark 4.1.3 and 4.2.0 worker sources clear the registry per task and call send_accumulator_updates; PythonRunner.handleEndOfDataSection receives those updates. The new native path has no corresponding forwarding or cleanup. Reproduction: /tmp/comet-6130-b0dab-review-current-1rkq0_j5/reference.py, src/bin/accumulator.rs, and reference.log. This validates the native component and Spark reference, not a full Comet query.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for catching this. I confirmed that the embedded interpreter updates its local accumulator registry but does not forward task updates to Spark. In d6588a6, accumulator-bearing Arrow UDFs now fall back to Spark's ArrowEvalPythonExec worker path. I added an end-to-end test that verifies the plan, result rows, and driver accumulator value of 4 with Comet enabled. The PyArrow UDF jobs pass on Spark 4.1 and 4.2; the remaining Rust and Spark SQL checks are still running.

@viirya

viirya commented Oct 6, 2026

Copy link
Copy Markdown
Member Author

I also addressed the exchange-identity issue. Equality and hashing no longer compare the full Arrow UDF protobuf, which contains generated argument expression IDs. They compare the Catalyst UDF expressions that Spark canonicalizes, while retaining the batch limits in plan identity. The new regression checks sameResult for equivalent UDF plans with different generated IDs, and checks that different functions and batch limits remain distinct.

These changes are in b0dab2f53. Local Spark 4.1 and 4.2 test compilation, the plan-equality suite, and Python syntax checks passed. The exact-head CI run passed: 34 successful checks and 14 skipped, including the Spark 4.1 and 4.2 native PyArrow UDF suites.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Scalar Arrow UDFs use Spark’s Python workers, interrupting Comet’s native pipeline.
  • Design approach: Add an opt-in PyO3 operator with Arrow C Data transfers.
  • Correctness / compatibility analysis: Compared Spark 4.1.3 and 4.2.0 sources and checked the earlier-version shims. Previously reported cases are addressed. Found one additional P2 issue: distributed archives bypass the new file fallback.
  • Key design decisions: Feature gating, version shims, and the scalar-type allow-list contain the scope. The operator and Python bridge have clear responsibilities without unnecessary abstraction.
  • Implementation sketch: Scala serializes callables, arguments, types, and batch limits. Rust creates per-partition callables, slices inputs, validates Python results, and appends output columns.
  • Behavioral changes worth calling out: Compared the touched paths with branch-1.1 at 7b7eec69282a7abc33a0e8e687ca594c72f140ec. Native scalar Arrow execution is an intended opt-in addition. The archive failure is unintended. Shared interpreter state, GIL contention, and Python memory outside Comet’s accounting are documented. Reported benchmarks distinguish PyArrow-kernel gains from GIL-bound slowdowns.
  • Suggested improvements: Include listArchives() in the distributed-file fallback and add an archive regression test.

Reviewed the entire 32-file diff at d6588a6c4c935d43a77211496044782df8673d88 against 99c05eb4f2a3facc2476f97be4aa24a4a466c8ea. Confirmed non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 38 successful checks, 14 skipped, none failed or pending. Linux builds, Rust tests, Spark 4.1 SQL suites, and Spark 4.0–4.2 PyArrow jobs passed. Logs confirm eight feature Rust tests and four native Scala cases on each of Spark 4.1 and 4.2, including the accumulator fallback regression.

Validation: Eight unit tests and two focused operator checks passed locally in an isolated harness importing the unmodified head sources. Spark 4.1.3 reference queries verified batch caps and accumulator updates. A separate archive reproduction demonstrated the finding below. Local Python setup issues were corrected before the successful harness run. No full local Comet JNI/JVM build, local Spark 4.2 execution, macOS validation, or benchmark rerun was performed.

if (op.conf.pythonUDFProfiler.nonEmpty) {
return Unsupported(Some("Arrow UDF profiling is not supported in-process"))
}
if (SparkSession.active.sparkContext.listFiles().nonEmpty) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could this guard also check SparkContext.listArchives()? Spark stores addArchive() resources separately, so listFiles() remains empty. With native Arrow UDF execution enabled, a scalar UDF that reads an extracted archive through SparkFiles.get() therefore passes the support checks, although embedded Python has no Spark file setup. A ZIP containing offset.txt with value 7, read by a long-valued UDF over [1, 2, 3, 4], returns [8, 9, 10, 11] in Spark but fails with AssertionError in the native operator. This breaks valid UDFs using distributed model or data archives. Keep archive-bearing contexts on Spark’s worker path until that setup is supported.

Evidence: Spark 4.1.3 reproduction reported listFiles=0, listArchives=1, pythonIncludes=0, broadcastVars=0, environment={PYTHONHASHSEED=0}, and no accumulator marker in the real PySpark command. The query returned [8, 9, 10, 11]. An isolated harness importing this head’s unmodified ArrowPythonUdfExec and ArrowPythonUdf evaluated the same serialized command and returned Arrow UDF Python error: AssertionError. Spark 4.1.3 and 4.2.0 SparkContext sources confirm separate file/archive registries. Reproduction artifacts: /tmp/comet-6130-review-1791283631/archive_reference.py, archive_reference.log, harness/tests/operator.rs, and archive-native.log. This validates the Spark reference and native component, not a full Comet query.

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

Labels

area:udf enhancement New feature or request run-pyarrow-udf-tests Run the PyArrow UDF tests on this pull request instead of waiting for the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support scalar PySpark Arrow UDFs in Comet native execution via PyO3

4 participants