Repository navigation
Conversation
| 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" |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| "_import_from_c", | ||
| ( | ||
| &raw const ffi_array as Py_uintptr_t, | ||
| &raw const ffi_schema as Py_uintptr_t, |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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?) |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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() |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
|
Addressed both review comments in |
sunchao
left a comment
There was a problem hiding this comment.
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) { |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
|
This is a light fully automated review since there are so many PRs open. Unlike its sibling operators,
|
|
Thanks for catching both issues. In I also documented that embedded Python and PyArrow allocations are excluded from Rust's |
sunchao
left a comment
There was a problem hiding this comment.
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], |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
| case udf if udf.func.pythonIncludes != null && !udf.func.pythonIncludes.isEmpty => | ||
| "Arrow UDF Python includes are not supported in-process" |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
|
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 These changes are in |
sunchao
left a comment
There was a problem hiding this comment.
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.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec. 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) { |
There was a problem hiding this comment.
[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.
Which issue does this PR close?
Closes #6129.
Rationale for this change
Starting with Spark 4.1, scalar
@arrow_udfexpressions fall back atArrowEvalPythonExec, 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?
python-udfCargo feature and a nativeArrowPythonUdfExecthat invokes Python through PyO3 and exchanges arrays through the Arrow C Data interface. Python evaluation usesblock_in_place, preserving synchronous polling for JVM scans while allowing Tokio to hand off other work on native-source paths.TimeType, so the native path preserves that behavior. Profiler-enabled UDFs also fall back. Earlier Spark versions retain their current behavior.spark.comet.exec.nativeArrowPythonUDF.enabled, defaulting tofalse, plus documentation of embedded Python setup and process-level limitations.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 warningspassed; five Rust unit tests passed with the feature enabled.pyspark.cloudpickleand 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.cargo fmt --check, and CI configuration checks passed. Spark 3.5 and 4.0 compatibility profiles compiled earlier in this PR.spawn_blockingwithblock_in_place, a Spark 4.1.3 benchmark of 50 million rows withpyarrow.compute.negateand 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.local[16]: Spark median 0.143 s and native median 0.492 s (native 3.44x slower). The priorspawn_blockingversion 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-testsandrun-pyarrow-udf-testslabels request broader CI coverage before merge.