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]
