dongjoon-hyun commented on PR #58978:
URL: https://github.com/apache/spark/pull/58978#issuecomment-5805443736

   Thank you for this proposal, @viirya. In-process execution via JEP and Arrow 
CDI is an interesting direction for Arrow UDFs. I appreciate the benchmark 
numbers and the tests covering the cancellation and CDI resource lifecycle.
   
   I left 15 inline comments in [this 
review](https://github.com/apache/spark/pull/58978#pullrequestreview-5298451065).
 Here is a summary, roughly by priority.
   
   **Correctness (wrong results or crashes)**
   - Non-root `LIMIT`/`OFFSET` after `ORDER BY` can return arbitrary rows. 
`InProcessArrowEvalExec` is not an `EvalPythonExec`, so 
`InsertSortForLimitAndOffset` does not insert the local sort.
   - A map result whose entries child is sliced is not normalized by 
`_has_offset`. Arrow Java 19's importer ignores `ArrowArray.offset`, so the JVM 
silently reads shifted data.
   - Map key/value field names (and field metadata) are not validated. A map 
with non-standard child names passes validation and then fails with an NPE in 
`ArrowColumnVector`.
   - `result.cast(expected_type)` rejects valid struct results when a 
non-nullable child is null only under null parents.
   - The Python linter (mypy `disallow_untyped_defs`) fails on the new 
`pyspark.inprocess` modules.
   
   **Lifecycle and concurrency**
   - `invoke` after plugin `shutdown()` hits an NPE, and `register` recreates 
an interpreter without `sitePackages`.
   - `close()` always calls `release()` under the global, non-cancellable lock, 
even for tasks that registered nothing.
   - `shutdown()` can block executor/SparkContext stop behind a running Python 
batch.
   
   **Python API contract**
   - For primitive return types, the wrapper silently casts strings, 
timestamps, dates and similar values.
   - Zero-argument functions are accepted but fail at runtime.
   - The function is pickled eagerly at decoration time, so later globals are 
missing or stale.
   - `spark.udf.register(name, wrapper)` silently registers a `StringType` 
`SQL_BATCHED_UDF`.
   - `sitePackages` entries are appended with `sys.path.extend`: `.pth` files 
are not processed, and system site-packages take precedence.
   
   **Performance and design**
   - The whole closure is converted byte by byte on every task, on the single 
interpreter thread.
   - The new exec duplicates `EvalPythonExec`/`EvalPythonEvaluatorFactory` and 
has already drifted from it (no metrics, no named-argument unwrap, missed by 
physical rules keyed on `EvalPythonExec`). An in-process evaluator factory 
inside `ArrowEvalPythonExec` might address several of the points above at once.
   
   Separately, the docs page contradicts itself: "Must satisfy 
`spark.executor.cores == spark.task.cpus`" vs. "not a correctness requirement", 
and "`MapType` is not currently supported" vs. the map handling in the runtime.
   
   Thank you again for working on this.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to