dongjoon-hyun commented on code in PR #58978:
URL: https://github.com/apache/spark/pull/58978#discussion_r4110615659


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,308 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.io.File
+import java.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, PyConfig, 
SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.{PythonException, PythonUtils}
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  val SITE_PACKAGES_CONFIG = "spark.inprocess.python.sitePackages"
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var configured = false
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  private def configureInterpreter(sitePackages: Seq[String]): Unit = {
+    if (!configured) {
+      // Like Python workers, use a stable default hash seed on every 
executor. This must
+      // happen before JEP creates its process-wide main interpreter, 
including on restarts.
+      MainInterpreter.setInitParams(new 
PyConfig().setHashSeed(0).setUseHashSeed(true))
+      // JEP imports its Python package during construction, before our 
bootstrap runs.
+      SharedInterpreter.setConfig(new 
JepConfig().addIncludePaths(sitePackages: _*))
+      configured = true
+    }
+  }
+
+  private[python] def bootstrapScript(script: String): String = {
+    "try:\n" + script.linesIterator.map("    " + _).mkString("\n") +
+      "\nexcept BaseException as _bootstrap_error:\n" +
+      "    raise RuntimeError('In-process Python bootstrap failed: ' + " +
+      "repr(_bootstrap_error)) from None\n"
+  }
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      active.requireCompatible(sitePackages)
+    } else {
+      configureInterpreter(sitePackages)
+      val candidate = new InterpreterSession(sitePackages)
+      try {
+        candidate.initialize()
+        active = candidate
+      } catch {
+        case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
candidate.shutdown() }
+      }
+    }
+  }
+
+  def currentSession: InterpreterSession = synchronized {
+    checkState(active != null && active.isRunning)
+    active
+  }
+
+  def shutdown(): Unit = {
+    val session = synchronized { active }
+    if (session != null) session.shutdown()
+  }
+
+  private def checkState(running: Boolean): Unit = {
+    checkState(running, "In-process Python is not running; initialize the 
executor plugin first")
+  }
+
+  private def checkState(running: Boolean, message: String): Unit = {
+    if (!running) throw new IllegalStateException(message)
+  }
+
+  /**
+   * Tasks retain this generation, so stale tasks cannot enter a later 
SparkContext's interpreter.
+   * Lifecycle operations only hold the monitor while enqueueing work, never 
while running Python.
+   */
+  private[python] class InterpreterSession(val sitePackages: Seq[String] = 
Seq.empty) {
+    // CPython native calls need more stack than the usual JVM thread default. 
This is a
+    // platform-dependent size request, not protection against arbitrary 
native crashes.
+    private val executor = Executors.newSingleThreadExecutor(new ThreadFactory 
{
+      override def newThread(runnable: Runnable): Thread = {
+        val thread = new Thread(null, runnable, "inprocess-python", 8L * 1024 
* 1024)
+        thread.setDaemon(true)
+        thread
+      }
+    })
+    @volatile private var running = true
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+
+    def isRunning: Boolean = running
+    def isTerminated: Boolean = executor.isTerminated
+
+    def requireCompatible(paths: Seq[String]): Unit = {
+      if (!isRunning) {
+        throw new LifecycleException("In-process Python is still stopping. 
Wait for outstanding " +
+          "native work to finish or replace the executor process before 
starting a new context.")
+      }
+      if (sitePackages != paths) {
+        throw new LifecycleException("In-process Python is already running 
with different " +
+          "sitePackages. Stop the existing context before changing interpreter 
configuration.")
+      }
+    }
+
+    private[python] def onInterpreterThread[T](body: => T): T = {
+      val context = Option(TaskContext.get())
+      context.foreach(_.killTaskIfInterrupted())
+      val gate = new Object
+      var started = false
+      var cancelled = false
+      val future = synchronized {
+        checkState(running)
+        executor.submit(new Callable[T] {
+          override def call(): T = {
+            gate.synchronized {
+              if (cancelled) throw new TaskKilledException("Cancelled before 
Python invocation")
+              started = true
+            }
+            body
+          }
+        })
+      }
+      var interrupted = false
+      try {
+        while (true) {
+          val taskCancelled = context.exists(_.isInterrupted())
+          if (interrupted || taskCancelled) {
+            val cancelledBeforeStart = gate.synchronized {
+              if (started) false else {
+                cancelled = true
+                future.cancel(false)
+                true
+              }
+            }
+            if (cancelledBeforeStart) {
+              context.foreach(_.killTaskIfInterrupted())
+              throw new InterruptedException("Cancelled before Python 
invocation")
+            }
+          }
+          try {
+            val result = future.get(100, TimeUnit.MILLISECONDS)
+            context.foreach(_.killTaskIfInterrupted())
+            return result
+          } catch {
+            case _: TimeoutException =>
+            case _: InterruptedException => interrupted = true
+            case e: ExecutionException => throw e.getCause
+          }
+        }
+        throw new IllegalStateException("Unreachable")
+      } finally {
+        // Once native work starts, wait for it even after cancellation: the 
caller still owns
+        // CDI structs that Python may use. Pending work, however, is safe to 
cancel immediately.
+        if (interrupted) Thread.currentThread().interrupt()
+      }
+    }
+
+    def initialize(): Unit = onInterpreterThread {
+      val candidate = new SharedInterpreter()
+      try {
+        candidate.set("_site_packages", sitePackages.asJava)
+        val sparkPaths = 
PythonUtils.sparkPythonPath.split(File.pathSeparator).filter(_.nonEmpty)
+        candidate.set("_spark_paths", sparkPaths.toSeq.asJava)
+        candidate.exec(bootstrapScript(
+          """import os, site, sys
+            |_configured = [os.path.abspath(p) for p in _site_packages]
+            |_before = set(sys.path)
+            |for _path in _configured:
+            |    site.addsitedir(_path)
+            |_added = [p for p in sys.path if p not in _before and p not in 
_configured]
+            |_preferred = list(dict.fromkeys(list(_spark_paths) + _configured 
+ _added))

Review Comment:
   **[High] On YARN, a PySpark installed in the venv shadows Spark's own 
`pyspark.zip`.**
   
   `_spark_paths` comes from `PythonUtils.sparkPythonPath`, which includes 
`$SPARK_HOME/python/lib/pyspark.zip` and the py4j zip only when `SPARK_HOME` is 
set. YARN containers don't get `SPARK_HOME` 
(`ExecutorRunnable.prepareEnvironment`). There, `pyspark.zip`/py4j arrive only 
through the container `PYTHONPATH` (`Client.scala` 
`setExecutorEnv("PYTHONPATH", ...)`), so the only `_spark_paths` entry is the 
spark-core jar, which has no Python sources. This reordering then puts 
`_configured` + `_added` ahead of the env `PYTHONPATH` entries.
   
   Scenario: the documented YARN recipe (`--archives myvenv.zip#myvenv`, 
`spark.inprocess.python.sitePackages=./myvenv/...`), where the venv has a PyPI 
`pyspark`, e.g. pulled in by `delta-spark`. `from pyspark.inprocess.runtime 
import ...` resolves to the venv copy. A released PySpark has no 
`pyspark.inprocess`, so bootstrap fails with `ModuleNotFoundError`, plugin init 
rethrows, and every executor fails to start. A matching future release would 
still produce driver/executor PySpark version skew.
   
   The worker path avoids this: `PythonWorkerFactory` merges `sparkPythonPath` 
with the env `PYTHONPATH` ahead of site-packages. The guide's statement that 
Spark distribution paths take precedence over these directories doesn't hold on 
YARN, and `test_spark_python_distribution_precedes_site_packages` only covers 
local mode with `SPARK_HOME` set. Suggestion: also promote the process 
`PYTHONPATH` entries, the same way `PythonWorkerFactory` does.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,322 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+
+"""Arrow CDI entry points called on the executor's dedicated JEP interpreter 
thread.
+
+Functions are registered once per task and released when that task finishes. 
Calls
+pass only a handle and CDI addresses, so large closures are not copied per 
batch.
+"""
+
+import sys
+from typing import Any, Callable, Iterable, Optional, Sequence
+
+import pyarrow as pa

Review Comment:
   **[Low] The in-process path never checks the minimum PyArrow version.**
   
   `arrow_udf` / `pandas_udf` on the driver and the worker's Arrow handler all 
call `require_minimum_pyarrow_version()` (minimum 18.0.0). Neither 
`inprocess_udf` nor this runtime calls it, and nothing imported here enforces 
it indirectly. The 18.0.0+ requirement is only stated in the guide's table.
   
   With an older pyarrow on the driver or executors, users would get whatever 
`TypeError` / `AttributeError` a version-specific API or behavior produces 
during registration or validation, instead of a clear 
`UNSUPPORTED_PACKAGE_VERSION`. I have not reproduced this against pyarrow 
15-17. Suggestion: call `require_minimum_pyarrow_version()` in 
`InProcessUDFWrapper` and in `_inprocess_register`.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalPythonExec.scala:
##########
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import org.apache.spark.sql.catalyst.expressions.{Attribute, PythonUDF}
+import org.apache.spark.sql.execution.SparkPlan
+
+/** Row-based CDI execution, sharing the standard Python UDF evaluator 
contracts. */
+case class InProcessArrowEvalPythonExec(
+    udfs: Seq[PythonUDF],
+    resultAttrs: Seq[Attribute],
+    child: SparkPlan) extends EvalPythonExec with PythonSQLMetrics {
+
+  override protected def evaluatorFactory: EvalPythonEvaluatorFactory = {
+    InProcessPythonUDFBuilder.checkWorkerEnvironment(conf)

Review Comment:
   **[Medium] Worker-only settings are silently ignored.**
   
   Only `spark.pythonWorkerEnv.*` is rejected here. The in-process path also 
ignores the following, without rejecting them or mentioning them in the guide:
   
   - `spark.executor.pyspark.memory`: the worker enforces it with `RLIMIT_AS` 
via `PYSPARK_EXECUTOR_MEMORY_MB` / `setup_memory_limits`. In-process, a UDF 
that would have hit a Python `MemoryError` instead grows the executor JVM's 
native memory until the container is OOM-killed, losing all of that executor's 
tasks. `ResourceProfile` still adds this amount to the container request, but 
nothing enforces it.
   - `spark.sql.pyspark.udf.profiler` (`perf`/`memory`): `spark.profile.show()` 
reports nothing for in-process UDFs, with no warning. The profiler results also 
rely on accumulators, which this path rejects.
   
   Suggestion: either reject these settings the same way as `pythonWorkerEnv`, 
or document them explicitly as unsupported.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,308 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.io.File
+import java.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, PyConfig, 
SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.{PythonException, PythonUtils}
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  val SITE_PACKAGES_CONFIG = "spark.inprocess.python.sitePackages"
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var configured = false
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  private def configureInterpreter(sitePackages: Seq[String]): Unit = {
+    if (!configured) {
+      // Like Python workers, use a stable default hash seed on every 
executor. This must
+      // happen before JEP creates its process-wide main interpreter, 
including on restarts.
+      MainInterpreter.setInitParams(new 
PyConfig().setHashSeed(0).setUseHashSeed(true))
+      // JEP imports its Python package during construction, before our 
bootstrap runs.
+      SharedInterpreter.setConfig(new 
JepConfig().addIncludePaths(sitePackages: _*))
+      configured = true

Review Comment:
   **[Medium] `configured = true` is set before JEP actually applies the 
config.**
   
   `SharedInterpreter.setConfig` only replaces a static field. JEP applies the 
config, and sets its own `initialized = true`, only after a `SharedInterpreter` 
is constructed successfully (the `import jep` inside `configureInterpreter` 
must succeed).
   
   Scenario: local mode or a notebook that reuses one JVM.
   1. The first `SparkContext` has a typo in 
`spark.inprocess.python.sitePackages`, which per the guide is where `jep` must 
be importable from.
   2. `new SharedInterpreter()` fails, but `configured` is already `true`.
   3. The user stops the session, fixes the conf and starts again. 
`configureInterpreter` is skipped, the new `SharedInterpreter` uses the stale 
`JepConfig`, and it fails the same way.
   4. This repeats until the process is restarted.
   
   `setInitParams` can't simply be re-run: it throws `IllegalStateException` 
once the main interpreter exists. So `setInitParams` and `setConfig` need 
separate tracking, and `setConfig` should be retried until an interpreter has 
been created successfully.
   
   Minor: `val candidate = new SharedInterpreter()` (L179) sits outside the 
`try`, so nothing closes it if the constructor fails part-way.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,308 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.io.File
+import java.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, PyConfig, 
SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.{PythonException, PythonUtils}
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  val SITE_PACKAGES_CONFIG = "spark.inprocess.python.sitePackages"
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var configured = false
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  private def configureInterpreter(sitePackages: Seq[String]): Unit = {
+    if (!configured) {
+      // Like Python workers, use a stable default hash seed on every 
executor. This must
+      // happen before JEP creates its process-wide main interpreter, 
including on restarts.
+      MainInterpreter.setInitParams(new 
PyConfig().setHashSeed(0).setUseHashSeed(true))

Review Comment:
   **[High] Non-isolated `PyConfig` lets CPython's `faulthandler` replace 
HotSpot's fatal-signal handlers.**
   
   `new PyConfig()` is JEP's non-isolated config (`PyConfig_InitPythonConfig` 
in `pyembed.c`), and JEP does not forward any `faulthandler` / `dev_mode` / 
`install_signal_handlers` control. So CPython honors `PYTHONFAULTHANDLER` / 
`PYTHONDEVMODE` from the executor environment and `_PyFaulthandler_Init` 
installs `sigaction` handlers for `SIGSEGV`/`SIGBUS`/`SIGFPE`/`SIGILL` over 
HotSpot's.
   
   Scenario: `spark.executorEnv.PYTHONFAULTHANDLER=1` (a common PySpark 
debugging tip) or a Python base image that sets it. HotSpot raises `SIGSEGV` 
routinely (safepoint polls, implicit null checks). On the next one, 
`faulthandler_fatal_error` prints `Fatal Python error: Segmentation fault`, 
restores the previous handler and re-raises from libc `raise()`. HotSpot cannot 
map that pc to compiled code, so it aborts with an hs_err file and the executor 
dies shortly after plugin init.
   
   Suggestion: use `PyConfig.isolated()` / `setUseEnvironment(false)` and pass 
the required paths explicitly, or reject/unset these variables (e.g. a 
`System.getenv` check) before initializing the interpreter.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonPlugin.scala:
##########
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.util.{Map => JMap}
+
+import scala.util.control.NonFatal
+
+import org.apache.spark.api.plugin.{DriverPlugin, ExecutorPlugin, 
PluginContext, SparkPlugin}
+import org.apache.spark.internal.Logging
+
+/**
+ * Spark plugin that initializes jep's SharedInterpreter on a dedicated 
executor thread,
+ * enabling in-process Python UDF execution with zero-copy Arrow data passing.
+ *
+ * Register via Spark config:
+ *   spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin
+ *
+ * Requirements:
+ *  - jep (Java Embedded Python) must be on the executor classpath (provided 
scope)
+ *  - Python 3.11+ with PyArrow 18+ and PySpark installed in the executor 
environment
+ *
+ * Calls from concurrent tasks are serialized on the interpreter thread. One 
task per executor
+ * is recommended for throughput but is not required for correctness.
+ *
+ * @see [[InProcessPythonRuntime]] for the interpreter singleton
+ */
+class InProcessPythonPlugin extends SparkPlugin {
+  override def driverPlugin(): DriverPlugin = null
+
+  override def executorPlugin(): ExecutorPlugin = new 
InProcessPythonExecutorPlugin()
+}
+
+private[python] class InProcessPythonExecutorPlugin extends ExecutorPlugin 
with Logging {
+
+  override def init(ctx: PluginContext, extraConf: JMap[String, String]): Unit 
= {
+    logInfo("Initializing in-process Python runtime (jep SharedInterpreter).")
+    try {
+      val sitePackages = ctx.conf()
+        .getOption(InProcessPythonRuntime.SITE_PACKAGES_CONFIG)
+        .map(_.split(",").map(_.trim).filter(_.nonEmpty).toSeq)
+        .getOrElse(Seq.empty)
+      InProcessPythonRuntime.initialize(sitePackages)

Review Comment:
   **[Medium] Plugin init never checks for the `provided` `arrow-c-data` JAR.**
   
   Init exercises JEP and the Python bootstrap. `org.apache.arrow.c.*` is first 
resolved only inside a task (`ArrowArray.allocateNew` in the evaluator, and 
`InProcessArrowBridge`).
   
   Scenario: the executor has `jep.jar` but no `arrow-c-data` JAR. The plugin 
logs "initialized successfully" and executors register. Then every in-process 
UDF task fails with `NoClassDefFoundError: org/apache/arrow/c/ArrowArray` and 
is retried until `spark.task.maxFailures`, far from the cause. The init error 
hint also mentions only `jep.jar`.
   
   Suggestion: resolve `org.apache.arrow.c.Data` here, e.g. with 
`Class.forName`. Optionally also run a tiny export/import round trip so the CDI 
JNI library gets loaded as well, and fail fast with a hint that names the CDI 
JAR.



##########
docs/sql-pyspark-inprocess-udf.md:
##########
@@ -0,0 +1,640 @@
+---
+layout: global
+title: In-Process Python UDFs
+displayTitle: In-Process Python UDFs
+license: |
+  Licensed to the Apache Software Foundation (ASF) under one or more
+  contributor license agreements.  See the NOTICE file distributed with
+  this work for additional information regarding copyright ownership.
+  The ASF licenses this file to You under the Apache License, Version 2.0
+  (the "License"); you may not use this file except in compliance with
+  the License.  You may obtain a copy of the License at
+
+     http://www.apache.org/licenses/LICENSE-2.0
+
+  Unless required by applicable law or agreed to in writing, software
+  distributed under the License is distributed on an "AS IS" BASIS,
+  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+  See the License for the specific language governing permissions and
+  limitations under the License.
+---
+
+* Table of contents
+{:toc}
+
+## Runtime and result contract
+
+Each executor owns a dedicated interpreter thread. The plugin initializes the
+interpreter on that thread, and task calls and shutdown are dispatched to the
+same thread. The JVM is asked to allocate an 8 MiB stack for this thread; the
+actual size is platform-dependent. Calls from concurrent tasks are queued on 
the
+interpreter thread.
+One task per executor is recommended for throughput, but is not a correctness 
requirement.
+Application-level Python parallelism comes from multiple executor JVMs.
+The plugin configures JEP's process-wide interpreter with hash seed `0`, 
matching
+Spark's default Python worker seed. It must initialize before any other JEP 
user in
+the JVM. The seed cannot change between SparkContexts in the same process; a 
custom
+worker `PYTHONHASHSEED` does not override this embedded-runtime setting.
+
+Task cancellation cannot safely stop arbitrary native Python code. An 
interrupted
+caller waits for the current invocation to finish before freeing the Arrow CDI
+structures, then restores its interrupt status. A UDF that never returns can
+therefore prevent its task from completing cancellation and block every 
subsequent
+in-process UDF on that executor, including calls from other tasks, jobs, and 
sessions.
+Recovery from a permanently hung invocation requires replacing the executor 
process.
+Plugin shutdown stops accepting new calls and waits up to five seconds for the 
interpreter thread. If a call is
+still running, cleanup stays queued behind it; its memory remains live until 
the
+call returns or the process exits. Shutdown does not forcibly interrupt native
+code. A new interpreter cannot start until the previous one has fully stopped.
+
+A scalar UDF must return a `pyarrow.Array` with exactly one element per input 
row.
+The runtime checks the result type against the declared Spark type, including
+nested fields, decimal scale, and timestamp unit/timezone. Value types must 
match
+exactly: use an explicit PyArrow cast in the UDF for numeric or other 
conversions.
+Nested field nullability may differ if the actual values satisfy the declared 
nullability. Sliced results, including nested
+child slices, are copied to remove offsets that Arrow Java's CDI importer 
cannot
+read. Compatible results retain zero-copy transfer.
+
+The API produces a regular `PythonUDF` expression with an in-process evaluation
+type. Spark's existing `ArrowEvalPython` planning rules handle aggregation,
+nested calls, nondeterminism, and filter/limit pushdown. A dedicated
+`InProcessArrowEvalPythonExec` extends `EvalPythonExec`, reusing its 
projection,
+row queue, result join, and partition-evaluator path. Ordinary Python UDFs 
continue
+to use Python workers.
+
+`maxRecordsPerBatch <= 0` means no row-count limit. The independent
+`spark.sql.execution.arrow.maxBytesPerBatch` limit still applies when positive.
+Only UDF arguments are converted to Arrow. Other columns stay in Spark rows,
+buffered in a spillable queue until the results are joined back. Duplicate 
nested
+field names in UDF arguments or declared results are rejected before Arrow Java
+reads their buffers.
+
+Each batch uses fresh input buffers. A Python function may retain an input 
array;
+later batches do not overwrite it. Retained arrays keep native memory alive, so
+functions should release them when no longer needed. JVM input vectors and 
result
+vectors are released on task completion, early termination and failure.
+
+UDF deserialization uses PySpark's bundled cloudpickle. Each task registers its
+own function instance once and passes a small handle for subsequent batches.
+Task completion queues release of the registered function and its closure 
state. Imported
+Python modules still share executor-wide state. Configured site-packages paths
+are supplied to JEP before its first construction, so JEP itself can be found 
in
+an archived venv. They are then processed with `site.addsitedir`, including 
`.pth`
+files. Spark's own PySpark and Py4J distribution paths come first, followed by
+configured directories, newly discovered `.pth` paths, and existing system 
paths.
+Already imported modules cannot be replaced by changing the search path.
+
+Spark broadcasts, accumulators, `SparkContext.addPyFile`, and Python 
`TaskContext`
+are not supported by this embedded runtime. Captured broadcast and accumulator
+objects are rejected during serialization; functions must not access them 
through
+imported modules either. Install modules on executors before startup, 
optionally
+using `spark.inprocess.python.sitePackages`. Session-scoped 
`spark.pythonWorkerEnv.*`
+settings are rejected: the shared interpreter cannot apply per-session process
+environments. Configure environment variables before the executor starts, for 
example
+with `spark.executorEnv.NAME` (or the launching environment in local mode).
+SQL registration through `spark.udf.register` is not supported and is rejected 
at registration time.
+Spark Connect does not support this execution mode; both client SQL 
registration and
+server planning reject it. The decorator accepts a `DataType` or a DDL string; 
DDL
+strings are parsed lazily with the active Spark session. It exposes `func`, 
`returnType`,
+`evalType`, `deterministic`, and `asNondeterministic()` along with the 
function's name
+and docstring.
+Functions must receive at least one input column (a literal also works) to 
determine
+the batch length. Positional and keyword arguments are supported. Functions are
+serialized on first use, so globals can be defined or rebound after decoration
+and before that first call. The driver's Python major.minor
+version must match the embedded interpreter; registration checks this before
+unpickling. Python exceptions, including `SystemExit` during deserialization or
+execution, are converted into task failures. Tracebacks honor the query's
+`spark.sql.execution.pyspark.udf.hideTraceback.enabled`,
+`spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled`, and
+`spark.sql.execution.pyspark.udf.tracebackWithLocals.enabled` settings. Native 
process
+termination remains outside this exception handling.
+
+## Overview
+
+In-process Python UDFs embed CPython directly into the Spark executor JVM using
+[jep (Java Embedded Python)](https://github.com/ninia/jep), eliminating the 
IPC overhead of
+standard Python UDFs and pandas UDFs. Data is passed to Python as
+[PyArrow](https://arrow.apache.org/docs/python/) arrays via the
+[Arrow C Data 
Interface](https://arrow.apache.org/docs/format/CDataInterface.html) — zero-copy
+for compatible input and output buffers. Row-to-Arrow conversion and 
normalization
+of sliced results still copy data.
+
+**Use `inprocess_udf` when:**
+- You are already using `pandas_udf` for vectorized transformations and want 
lower latency.
+- Your UDF operates on Arrow/PyArrow arrays (e.g. using `pyarrow.compute`).
+- You can deploy enough executor JVMs for Python parallelism (see 
[Requirements](#requirements)).
+
+**Stick with `pandas_udf` or `udf` when:**
+- You need pandas Series semantics in your UDF logic.
+- You need concurrent Python invocations within a single executor.
+- You are not able to install jep on executors.
+
+---
+
+## Quick Start
+
+### 1. Install dependencies
+
+```bash
+pip install "jep>=4.3.2" pyarrow cloudpickle
+```
+
+JEP and `org.apache.arrow:arrow-c-data` are provided dependencies and are not
+bundled with Spark. Supply their JARs on the driver/executor classpaths before
+starting Spark, and make the JEP native library available. Use an 
`arrow-c-data`
+version matching Spark's Arrow Java version. Installing the Python packages 
alone
+does not supply the Arrow Java CDI JAR.
+
+Building JEP from source requires a JDK, a C compiler, and development headers 
for
+the Python version being embedded (for example, `python3.12-dev` on Ubuntu with
+Python 3.12). These headers are build dependencies; running a prebuilt 
compatible
+JEP installation does not require the development package. The corresponding
+Python shared library must remain available at runtime.
+
+### 2. Register the plugin
+
+```python
+spark = SparkSession.builder \
+    .config("spark.plugins",
+            "org.apache.spark.sql.execution.python.InProcessPythonPlugin") \
+    .config("spark.executor.cores", "1") \
+    .config("spark.task.cpus", "1") \
+    .getOrCreate()
+```
+
+### 3. Write and call a UDF
+
+```python
+import pyarrow.compute as pc
+from pyspark.inprocess.udf import inprocess_udf
+from pyspark.sql.types import LongType
+
+@inprocess_udf(return_type=LongType())
+def double(x):
+    return pc.multiply(x, 2)
+
+df = spark.range(10)
+df.select(double(df["id"])).show()
+```
+
+The function receives a `pa.Array` for each input column and must return a 
`pa.Array`.
+
+---
+
+## Examples
+
+### String transformation
+
+```python
+import pyarrow.compute as pc
+from pyspark.inprocess.udf import inprocess_udf
+from pyspark.sql.types import StringType
+
+@inprocess_udf(return_type=StringType())
+def upper(s):
+    return pc.utf8_upper(s)
+
+df = spark.createDataFrame([("hello",), ("world",)], ["text"])
+df.select(upper(df["text"])).show()
+# +------------+
+# |upper(text) |
+# +------------+
+# |HELLO       |
+# |WORLD       |
+# +------------+
+```
+
+### Multi-column UDF
+
+A UDF receives one `pa.Array` argument per input column:
+
+```python
+import pyarrow.compute as pc
+from pyspark.inprocess.udf import inprocess_udf
+from pyspark.sql.types import DoubleType
+
+@inprocess_udf(return_type=DoubleType())
+def weighted_sum(x, y):
+    return pc.add(pc.multiply(x, 0.6), pc.multiply(y, 0.4))
+
+df = spark.createDataFrame([(1.0, 2.0), (3.0, 4.0)], ["x", "y"])
+df.select(weighted_sum(df["x"], df["y"])).show()
+```
+
+### Closure capture
+
+Free variables are captured by cloudpickle and frozen into the serialized UDF. 
The captured
+value is captured at first use and shipped with the function to every executor.
+Rebinding a global before the first call is reflected in the serialized 
function;
+subsequent calls reuse the cached serialization:
+
+```python
+import pyarrow.compute as pc
+from pyspark.inprocess.udf import inprocess_udf
+from pyspark.sql.types import DoubleType
+
+SCALE_FACTOR = 100.0
+
+@inprocess_udf(return_type=DoubleType())
+def scale(x):
+    return pc.multiply(x, SCALE_FACTOR)
+```
+
+### Non-deterministic UDF
+
+Pass `deterministic=False` when the UDF produces different results for the 
same input (e.g.
+random sampling). This prevents the optimizer from deduplicating or reordering 
calls:
+
+```python
+import random
+import pyarrow as pa
+import pyarrow.compute as pc
+from pyspark.inprocess.udf import inprocess_udf
+from pyspark.sql.types import DoubleType
+
+@inprocess_udf(return_type=DoubleType(), deterministic=False)
+def add_noise(x):
+    noise = pa.array([random.gauss(0.0, 0.01) for _ in range(len(x))])
+    return pc.add(x, noise)
+```
+
+---
+
+## Requirements
+
+| Requirement | Detail |
+|---|---|
+| Python | 3.11+; driver and embedded major.minor versions must match |
+| jep | 4.3.2+ (`pip install jep`) |
+| `arrow-c-data` JAR | Provided separately; match Spark's Arrow Java version |
+| PyArrow | 18.0.0+ |
+| cloudpickle | Bundled with PySpark |
+| Python concurrency | One invocation at a time per executor (see below) |
+
+### Executor concurrency
+
+In-process UDFs use one `SharedInterpreter` on a dedicated thread per executor.
+Multiple Spark tasks can share an executor, including with fractional
+`spark.task.cpus`, but their Python invocations are serialized. `local[*]` 
therefore
+works but does not provide parallel embedded Python execution.
+
+For throughput, consider `spark.executor.cores=1, spark.task.cpus=1` and 
multiple
+executors. More executors also mean more JVM overhead; compare with 
worker-based
+Arrow UDFs under the same total CPU and memory budget.
+
+---
+
+## Deployment and Distribution
+
+### Local development
+
+For local development (e.g. `SparkSession.builder.master("local[*]")`), 
install jep and the
+required Python packages into the virtual environment you run PySpark from. 
The venv's
+site-packages must be supplied explicitly to the embedded interpreter. Use the
+PySpark distribution from the same Spark build; a separately pip-installed 
PySpark
+version may not contain this API or match the JVM classes.
+
+```bash
+python3 -m venv .venv
+.venv/bin/pip install "jep>=4.3.2" pyarrow
+source .venv/bin/activate
+```
+
+JEP and Arrow CDI must be on the JVM **system classpath** before the JVM 
starts.
+`--jars` alone only configures Spark's user classloader and is insufficient. 
The CDI
+JAR must match the Arrow Java version in the Spark build. For example:
+
+```bash
+JEP_DIR="$(python3 -c 'import importlib.util, pathlib; 
print(pathlib.Path(importlib.util.find_spec("jep").origin).parent)')"
+ARROW_C_DATA_JAR=/absolute/path/to/arrow-c-data.jar
+spark-submit --master 'local[1]' \
+  --driver-class-path "$JEP_DIR/*:$ARROW_C_DATA_JAR" \
+  --conf "spark.driver.extraLibraryPath=$JEP_DIR" \
+  --conf "spark.inprocess.python.sitePackages=$(dirname "$JEP_DIR")" \
+  --conf 
spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \
+  my_app.py
+```
+
+### Cluster deployment — prerequisite: build and zip the venv
+
+Both YARN and Kubernetes support distributing a virtual environment via 
`--archives`. Build the
+venv on a machine that matches the executor OS and Python version:
+
+```bash
+python3 -m venv myvenv
+myvenv/bin/pip install "jep>=4.3.2" pyarrow cloudpickle my-custom-lib
+(cd myvenv && zip -r ../myvenv.zip .)
+```
+
+Adjust `python3.11` in the paths below to match the Python version in your 
venv.
+
+---
+
+### YARN
+
+Spark extracts `--archives` to a relative path (`./myvenv/`) on each YARN 
container before the executor JVM
+starts. The key extra config compared to local development is
+`spark.executorEnv.PYSPARK_PYTHON`, which tells PySpark's Python worker to use 
the venv's
+Python executable (ensuring a consistent Python version between the 
JVM-embedded interpreter
+and any out-of-process fallbacks).
+
+```bash
+spark-submit \
+  --master yarn \
+  --deploy-mode cluster \
+  --archives myvenv.zip#myvenv \
+  --files /absolute/path/to/arrow-c-data.jar#arrow-c-data.jar \
+  --conf 
'spark.executor.extraClassPath=./myvenv/lib/python3.11/site-packages/jep/*:./arrow-c-data.jar'
 \
+  --conf 
spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \
+  --conf spark.executor.cores=1 \
+  --conf spark.task.cpus=1 \
+  --conf spark.executorEnv.PYSPARK_PYTHON=./myvenv/bin/python3 \
+  --conf 
spark.executor.extraJavaOptions="-Djava.library.path=./myvenv/lib/python3.11/site-packages/jep"
 \
+  --conf 
spark.inprocess.python.sitePackages=./myvenv/lib/python3.11/site-packages \
+  my_app.py
+```
+
+For local execution on the driver, use the driver classpath and native-library
+settings from the local example. A client-mode driver does not execute 
executor UDFs.
+
+---
+
+### Kubernetes
+
+#### Option A: Custom Docker image (recommended)
+
+Build an executor image from the same Spark distribution used by the driver. 
The
+image must contain Python 3.11 or newer, a matching Python shared library, and 
a
+JDK compatible with that Spark build. Do not assume `apache/spark:latest` has 
these
+versions or includes JEP's build dependencies.
+
+Install JEP and PyArrow into a known environment, such as `/opt/venv`, while
+building the image. Building JEP from its source distribution additionally 
requires
+a C compiler, the matching Python development headers, and a JDK (`JAVA_HOME` 
set).
+Then prepare the image with:
+
+- the JEP JAR and matching Arrow CDI JAR in `/opt/spark/jars`;
+- JEP's native library directory on the JVM library path before startup;
+- `spark.inprocess.python.sitePackages` pointing to `/opt/venv`'s 
site-packages;
+- Spark's matching `python/lib/pyspark.zip` and Py4J zip in the Spark 
distribution.
+
+The plugin adds Spark's Python distribution paths itself, ahead of any PySpark
+package installed in the venv. The submission example below assumes Python 3.11
+and `/opt/venv/lib/python3.11/site-packages/jep` for the native library 
directory;
+adjust both paths to the Python version used to build JEP.
+
+**Submit:**
+
+```bash
+spark-submit \
+  --master k8s://https://<k8s-api-server>:<port> \
+  --deploy-mode cluster \
+  --conf 
spark.kubernetes.container.image=my-registry/spark-inprocess:tested-build \
+  --conf 
spark.inprocess.python.sitePackages=/opt/venv/lib/python3.11/site-packages \
+  --conf 
spark.executor.extraLibraryPath=/opt/venv/lib/python3.11/site-packages/jep \

Review Comment:
   **[Medium] `spark.executor.extraLibraryPath` has no effect on Kubernetes.**
   
   Only Standalone (`StandaloneSchedulerBackend`) and YARN (`ExecutorRunnable`) 
read `EXECUTOR_LIBRARY_PATH`. Nothing under `resource-managers/kubernetes` 
applies it: `BasicExecutorFeatureStep` forwards only Java options and 
classpath, and `entrypoint.sh` does no `LD_LIBRARY_PATH` handling.
   
   If a user follows this example as written and the image doesn't set 
`LD_LIBRARY_PATH` itself, JEP's `System.loadLibrary("jep")` and its 
`LibraryLocator` fallback both fail. Plugin init then throws 
`UnsatisfiedLinkError` on every executor.
   
   Suggestion: in the K8s image recipe, set `LD_LIBRARY_PATH`, or use 
`-Djava.library.path` in `extraJavaOptions`. Alternatively, have the plugin 
call `MainInterpreter.setJepLibraryPath(...)` from `sitePackages`, since `jep/` 
must already be there.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,322 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+
+"""Arrow CDI entry points called on the executor's dedicated JEP interpreter 
thread.
+
+Functions are registered once per task and released when that task finishes. 
Calls
+pass only a handle and CDI addresses, so large closures are not copied per 
batch.
+"""
+
+import sys
+from typing import Any, Callable, Iterable, Optional, Sequence
+
+import pyarrow as pa
+import pyarrow.compute as pc
+
+from pyspark import cloudpickle
+from pyspark.errors import PySparkRuntimeError
+from pyspark.sql.pandas.types import to_arrow_type
+from pyspark.util import _format_exception
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+NullChecker = Callable[[pa.Array], None]
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType, NullChecker, 
bool, bool, bool]] = {}
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    timezone: str,
+    python_version: str,
+    large_var_types: bool = False,
+    hide_traceback: bool = False,
+    simplified_traceback: bool = False,
+    traceback_with_locals: bool = False,
+) -> None:
+    try:
+        embedded_version = "%d.%d" % sys.version_info[:2]
+        if python_version != embedded_version:
+            raise PySparkRuntimeError(
+                errorClass="PYTHON_VERSION_MISMATCH",
+                messageParameters={
+                    "worker_version": embedded_version,
+                    "driver_version": python_version,
+                },
+            )
+        # JEP exposes direct ByteBuffers through the buffer protocol. Unpickle 
a separate
+        # function per task without iterating over a PyJArray one JNI call per 
byte.
+        # Carry the type with the closure so driver-defined UDTs need no 
module import.
+        func, return_type = cloudpickle.loads(memoryview(serialized_udf))
+        expected_type = to_arrow_type(
+            return_type,
+            timezone=timezone,
+            prefers_large_types=large_var_types,
+            error_on_duplicated_field_names_in_struct=True,
+        )
+        if large_var_types:
+            expected_type = _large_binary_type(expected_type)
+        checker = _null_checker(expected_type) or (lambda array: None)
+        _udfs[handle] = (
+            func,
+            expected_type,
+            checker,
+            hide_traceback,
+            simplified_traceback,
+            traceback_with_locals,
+        )
+    except BaseException as error:
+        # In JEP, an uncaught SystemExit can terminate the entire executor JVM.
+        raise RuntimeError(
+            _UDF_TRACEBACK_SENTINEL
+            + _format_exception(error, hide_traceback, simplified_traceback, 
traceback_with_locals)
+        ) from None
+
+
+def _inprocess_release(handles: Iterable[str]) -> None:
+    for handle in handles:
+        _udfs.pop(handle, None)
+
+
+def _nullable_type(data_type: pa.DataType) -> pa.DataType:
+    def nullable_field(field: pa.Field) -> pa.Field:
+        return pa.field(field.name, _nullable_type(field.type), nullable=True)
+
+    if pa.types.is_struct(data_type):
+        return pa.struct([nullable_field(field) for field in data_type])
+    if pa.types.is_list(data_type):
+        return pa.list_(nullable_field(data_type.value_field))
+    if pa.types.is_large_list(data_type):
+        return pa.large_list(nullable_field(data_type.value_field))
+    if pa.types.is_map(data_type):
+        return pa.map_(
+            _nullable_type(data_type.key_type),
+            nullable_field(data_type.item_field),
+            keys_sorted=data_type.keys_sorted,
+        )
+    return data_type
+
+
+def _large_binary_type(data_type: pa.DataType) -> pa.DataType:
+    # The shared conversion keeps binary children for Variant and spatial 
types. Widen
+    # them only for CDI, where the layout must match the JVM, including nested 
occurrences.
+    if pa.types.is_binary(data_type):
+        return pa.large_binary()
+    if pa.types.is_struct(data_type):
+        return pa.struct([f.with_type(_large_binary_type(f.type)) for f in 
data_type])
+    if pa.types.is_list(data_type):
+        field = data_type.value_field
+        return pa.list_(field.with_type(_large_binary_type(field.type)))
+    if pa.types.is_map(data_type):
+        return pa.map_(
+            
data_type.key_field.with_type(_large_binary_type(data_type.key_type)),
+            
data_type.item_field.with_type(_large_binary_type(data_type.item_type)),
+            keys_sorted=data_type.keys_sorted,
+        )
+    return data_type
+
+
+# The predicate is deliberately conservative: hidden nulls may request a 
check, but a
+# null-free superset proves that all visible values satisfy the required-field 
contract.
+NullCheckPlan = tuple[Callable[[pa.Array], bool], NullChecker]
+
+
+def _null_check_plan(expected_type: pa.DataType) -> Optional[NullCheckPlan]:
+    def field_plan(field: pa.Field) -> Optional[NullCheckPlan]:
+        nested = _null_check_plan(field.type)
+        if field.nullable:
+            return nested
+
+        def needs_check(values: pa.Array) -> bool:
+            return bool(values.null_count) or (nested is not None and 
nested[0](values))
+
+        def check(values: pa.Array) -> None:
+            if values.null_count:
+                raise ValueError(
+                    f"In-process UDF returned nulls in non-nullable field 
{field.name}"
+                )
+            if nested is not None:
+                nested[1](values)
+
+        return needs_check, check
+
+    if pa.types.is_struct(expected_type):
+        fields = [(i, field_plan(f)) for i, f in enumerate(expected_type)]
+        checks = [(i, plan) for i, plan in fields if plan is not None]
+        if not checks:
+            return None
+
+        def needs_struct(array: pa.Array) -> bool:
+            return any(plan[0](array.field(i)) for i, plan in checks)
+
+        def check_struct(array: pa.Array) -> None:
+            valid = None
+            for i, (needs, check) in checks:
+                values = array.field(i)
+                if needs(values):
+                    if array.null_count:
+                        if valid is None:
+                            valid = pc.is_valid(array)
+                        # Filter only the child requiring a check, not its 
sibling payloads.
+                        values = pc.filter(values, valid)
+                    check(values)
+
+        return needs_struct, check_struct
+    if pa.types.is_list(expected_type) or 
pa.types.is_large_list(expected_type):
+        plan = field_plan(expected_type.value_field)
+        if plan is not None:
+
+            def check_list(array: pa.Array) -> None:
+                if plan[0](array.values):
+                    plan[1](pc.list_flatten(array))
+
+            return lambda array: plan[0](array.values), check_list
+    if pa.types.is_map(expected_type):
+        key_plan = _null_check_plan(expected_type.key_type)
+        item_plan = field_plan(expected_type.item_field)
+        # Arrow validation rejects null keys already; only their descendants 
need checks.
+        checks = [(i, p) for i, p in enumerate((key_plan, item_plan)) if p is 
not None]
+        if not checks:
+            return None
+
+        def entries(array: pa.Array) -> pa.Array:
+            start = array.offsets[0].as_py()
+            length = array.offsets[-1].as_py() - start
+            # values.field honors the entries struct's offset; keys/items do 
not.
+            return array.values.slice(start, length)
+
+        def needs_map(array: pa.Array) -> bool:
+            values = entries(array)
+            return any(plan[0](values.field(i)) for i, plan in checks)
+
+        def check_map(array: pa.Array) -> None:
+            if needs_map(array):
+                visible = pc.filter(array, pc.is_valid(array)) if 
array.null_count else array
+                values = entries(visible)
+                for i, (needs, check) in checks:
+                    if needs(values.field(i)):
+                        check(values.field(i))
+
+        return needs_map, check_map
+    return None
+
+
+def _null_checker(expected_type: pa.DataType) -> Optional[NullChecker]:
+    plan = _null_check_plan(expected_type)
+    return plan[1] if plan is not None else None
+
+
+def _has_offset(array: pa.Array) -> bool:
+    if array.offset:
+        return True
+    if pa.types.is_struct(array.type):
+        return any(_has_offset(array.field(i)) for i in 
range(array.type.num_fields))
+    if pa.types.is_list(array.type) or pa.types.is_large_list(array.type):
+        return _has_offset(array.values)
+    if pa.types.is_map(array.type):
+        return _has_offset(array.values)
+    return False
+
+
+def _with_schema(array: pa.Array, expected_type: pa.DataType) -> pa.Array:
+    # Rebind buffers after validating logical nullability. Arrow cast checks 
hidden child
+    # slots too, rejecting null children underneath null parents. from_buffers 
preserves
+    # those masks and applies the declared names, metadata and nullability 
without casting.
+    children = None
+    if pa.types.is_struct(expected_type):
+        children = [_with_schema(array.field(i), f.type) for i, f in 
enumerate(expected_type)]
+    elif pa.types.is_list(expected_type) or 
pa.types.is_large_list(expected_type):
+        children = [_with_schema(array.values, expected_type.value_type)]
+    elif pa.types.is_map(expected_type):
+        entries_type = pa.struct([expected_type.key_field, 
expected_type.item_field])
+        children = [_with_schema(array.values, entries_type)]
+    return pa.Array.from_buffers(
+        expected_type,
+        len(array),
+        array.buffers()[: array.type.num_buffers],
+        null_count=array.null_count,
+        children=children,
+    )
+
+
+def _validate_result(
+    result: pa.Array,
+    expected_rows: int,
+    expected_type: pa.DataType,
+    null_checker: Optional[NullChecker] = None,
+) -> pa.Array:
+    if not isinstance(result, pa.Array):
+        raise TypeError(f"In-process UDF must return a pyarrow.Array, got 
{type(result).__name__}")
+    if len(result) != expected_rows:
+        raise ValueError(f"In-process UDF returned {len(result)} rows; 
expected {expected_rows}")
+    if _nullable_type(result.type) != _nullable_type(expected_type):
+        raise TypeError(f"In-process UDF returned {result.type}; expected 
{expected_type}")
+    result.validate()
+    checker = null_checker if null_checker is not None else 
_null_checker(expected_type)
+    if checker is not None:
+        checker(result)
+    # Arrow Java's CDI importer does not honor ArrowArray.offset, including 
child offsets.
+    # Concatenation materializes the logical slice, preserving validity and 
nested values.
+    if _has_offset(result):
+        result = pa.concat_arrays([result])
+    return _with_schema(result, expected_type)
+
+
+def _inprocess_invoke(
+    handle: str,
+    input_array_ptrs: Sequence[int],
+    input_schema_ptrs: Sequence[int],
+    output_array_ptr: int,
+    output_schema_ptr: int,
+    expected_rows: int,
+    argument_names: Optional[Sequence[str]] = None,
+) -> None:
+    """Consume input CDI structs and export a validated, row-preserving result.
+
+    The caller owns the struct memory and releases unconsumed exports on 
failure.
+    Each batch owns its buffers; retained Python inputs are never overwritten.
+    """
+    hide_traceback = simplified_traceback = traceback_with_locals = False
+    try:
+        (
+            udf_func,
+            expected_type,
+            checker,
+            hide_traceback,
+            simplified_traceback,
+            traceback_with_locals,
+        ) = _udfs[handle]
+        if len(input_array_ptrs) != len(input_schema_ptrs):
+            raise ValueError("Mismatched input ArrowArray and ArrowSchema 
pointer counts")
+        input_arrays = [
+            pa.Array._import_from_c(int(ap), int(sp))
+            for ap, sp in zip(input_array_ptrs, input_schema_ptrs)
+        ]
+        names = argument_names if argument_names is not None else [""] * 
len(input_arrays)
+        if len(names) != len(input_arrays):
+            raise ValueError("Mismatched input argument names")
+        args = [value for name, value in zip(names, input_arrays) if not name]
+        kwargs = {str(name): value for name, value in zip(names, input_arrays) 
if name}
+        result = _validate_result(
+            udf_func(*args, **kwargs), int(expected_rows), expected_type, 
checker
+        )
+        result._export_to_c(int(output_array_ptr), int(output_schema_ptr))
+    except BaseException as error:
+        raise RuntimeError(

Review Comment:
   **[Medium] UDF errors can be lost or garbled when JEP converts the message 
to Java.**
   
   The traceback reaches the JVM only as the `JepException` message, which 
`withPythonException` searches for the sentinel. JEP 4.3.2 (`jep_exceptions.c`, 
`process_py_exception`) builds that message with `PyUnicode_AsUTF8(message)` 
and then JNI `NewStringUTF`, with no NULL check. `NewStringUTF` expects 
modified UTF-8.
   
   - Lone surrogates: e.g. `raise ValueError(os.fsdecode(b"/data/\xff.csv"))`, 
or a surrogate-escaped source path in the traceback. `PyUnicode_AsUTF8` returns 
NULL, so `getMessage()` is `null`. The user then gets `In-process Python 
infrastructure error: null`, and the exception type, message and traceback are 
all lost.
   - Non-BMP characters such as emoji come out as Latin-1 mojibake.
   - An embedded NUL truncates the message.
   
   Suggestion: make the text ASCII-safe before raising, e.g. `.encode("ascii", 
"backslashreplace").decode("ascii")`. Alternatively, return the formatted error 
as the function's return value instead of parsing the exception message.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonUDFBuilder.scala:
##########
@@ -0,0 +1,83 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.util.{Collections, List => JList}
+
+import scala.jdk.CollectionConverters._
+
+import org.apache.spark.api.python.{PythonEvalType, SimplePythonFunction}
+import org.apache.spark.sql.Column
+import org.apache.spark.sql.catalyst.expressions.PythonUDF
+import org.apache.spark.sql.catalyst.plans.logical.NamedParametersSupport
+import org.apache.spark.sql.classic.{ColumnNodeExpression, ExpressionUtils}
+import org.apache.spark.sql.internal.SQLConf
+import org.apache.spark.sql.types.DataType
+
+/**
+ * JVM-side builder for in-process [[PythonUDF]] expressions, called from the 
Python API
+ * via py4j's JVM reflection bridge 
(``sc._jvm.org.apache.spark...InProcessPythonUDFBuilder``).
+ *
+ * Accepts Java-typed arguments as passed by PySpark's ``sc._jvm`` proxy and 
returns a
+ * [[Column]] backed by a [[PythonUDF]] with the in-process evaluation type.
+ */
+object InProcessPythonUDFBuilder {
+
+  /**
+   * Build a [[Column]] backed by an in-process [[PythonUDF]] expression.
+   *
+   * @param name            display name (Python function ``__name__``)
+   * @param serializedFunc  cloudpickle bytes of the Python UDF
+   * @param returnTypeJson  JSON string of the Spark SQL return type
+   * @param jColumns        Java List of JVM [[Column]] objects (the UDF 
inputs)
+   * @param deterministic   whether the UDF always returns the same output for 
the same input;
+   *                        set to false for UDFs that use randomness or 
external state
+   * @param pythonVersion   driver's Python major.minor version
+   * @return                [[Column]] backed by an in-process [[PythonUDF]] 
expression
+   */
+  def build(
+      name: String,
+      serializedFunc: Array[Byte],
+      returnTypeJson: String,
+      jColumns: JList[Column],
+      deterministic: Boolean,
+      pythonVersion: String): Column = {
+    checkWorkerEnvironment(SQLConf.get)
+    val returnType = DataType.fromJson(returnTypeJson)
+    val inputExprs = jColumns.asScala.map(col => 
ColumnNodeExpression(col.node)).toSeq
+    NamedParametersSupport.splitAndCheckNamedArguments(inputExprs, name, 
SQLConf.get.resolver)
+    val function = new SimplePythonFunction(
+      serializedFunc,
+      Collections.emptyMap[String, String](),
+      Collections.emptyList[String](),
+      "",
+      pythonVersion,
+      Collections.emptyList(),
+      null)
+    ExpressionUtils.column(PythonUDF(
+      name, function, returnType, inputExprs,
+      PythonEvalType.SQL_SCALAR_ARROW_INPROCESS_UDF, deterministic))
+  }
+
+  private[python] def checkWorkerEnvironment(conf: SQLConf): Unit = {
+    require(PythonWorkerEnvironment.read(conf).isEmpty,

Review Comment:
   **[Low] User-facing validation throws an uncategorized 
`IllegalArgumentException`.**
   
   `require(...)` shows up as `IllegalArgumentException("requirement failed: 
...")`, with no error condition and no SQLSTATE. 
`common/utils/src/main/resources/error/README.md` says "You should not 
introduce new uncategorized errors". `PythonWorkerEnvironment` already has 
`INVALID_SPARK_CONFIG.*PYTHON_WORKER_ENV*` conditions, so a sibling condition 
would fit.
   
   There is also a per-task effect. The same check runs in 
`InProcessArrowEvalPythonExec.evaluatorFactory`, and with the default 
`spark.sql.execution.usePartitionEvaluator=false` that is evaluated inside the 
task closure (`EvalPythonExec.doExecute`). So if the conf is set after the 
Column is built, a deterministic configuration error becomes a task failure 
retried up to `spark.task.maxFailures`.
   
   The `require`s in `InProcessArrowBridge` (L89, L112) and in the evaluator 
factory (L68) guard internal invariants, so `SparkException.internalError` fits 
them better.



##########
python/pyspark/inprocess/udf.py:
##########
@@ -0,0 +1,225 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+"""
+Python API for in-process UDF registration.
+
+Usage::
+
+    import pyarrow.compute as pc
+    from pyspark.inprocess import inprocess_udf
+    from pyspark.sql.types import LongType
+
+    @inprocess_udf(return_type=LongType())
+    def double(x):
+        # x is a pa.Array; return a pa.Array
+        return pc.multiply(x, 2)
+
+    df.select(double(df.value)).show()
+"""
+
+import io
+import sys
+from functools import update_wrapper
+from inspect import getfullargspec
+from typing import Any, Callable, Optional, Union
+
+from pyspark import Accumulator, Broadcast, cloudpickle
+from pyspark.errors import PySparkTypeError, PySparkValueError
+from pyspark.sql.column import Column
+from pyspark.sql.types import DataType, _parse_datatype_string
+from pyspark.util import PythonEvalType
+
+
+class _InProcessPickler(cloudpickle.CloudPickler):
+    def reducer_override(self, obj: Any) -> Any:
+        if isinstance(obj, (Broadcast, Accumulator)):
+            raise TypeError("In-process UDFs do not support Spark broadcasts 
or accumulators")
+        return super().reducer_override(obj)
+
+
+def _serialize_udf(func: Callable, return_type: DataType) -> bytes:
+    buffer = io.BytesIO()
+    _InProcessPickler(buffer).dump((func, return_type))
+    return buffer.getvalue()
+
+
+class InProcessUDFWrapper:
+    """
+    Wraps a Python function as an in-process UDF.
+
+    Returned by ``@inprocess_udf``. Calling an instance with Spark ``Column``
+    arguments creates a ``Column`` expression backed by ``PythonUDF``
+    on the JVM side.
+    """
+
+    def __init__(
+        self, func: Callable, return_type: Union[DataType, str], 
deterministic: bool = True
+    ) -> None:
+        if not isinstance(return_type, (DataType, str)):
+            raise PySparkTypeError(
+                errorClass="NOT_EXPECTED_TYPE",
+                messageParameters={
+                    "expected_type": "DataType or str",
+                    "arg_name": "return_type",
+                    "arg_type": type(return_type).__name__,
+                },
+            )
+        self._return_type = return_type
+        self._parsed_return_type: Optional[DataType] = None
+        self.evalType = PythonEvalType.SQL_SCALAR_ARROW_INPROCESS_UDF
+        self._deterministic: bool = deterministic
+        self._name: str = getattr(func, "__name__", "inprocess_udf")
+
+        argspec = getfullargspec(func)
+        if not argspec.args and argspec.varargs is None and not 
argspec.kwonlyargs:

Review Comment:
   **[Low] The 0-arg check rejects functions that take only `**kwargs`.**
   
   This guard ignores `argspec.varkw`, so `inprocess_udf("long")(lambda **cols: 
pc.add(cols["a"], cols["b"]))` raises `INVALID_PANDAS_UDF` "0-arg 
inprocess_udfs are not supported". Keyword arguments otherwise work end to end: 
`__call__` builds a `namedArgumentExpression`, and `_inprocess_invoke` calls 
`udf_func(**kwargs)`. The guide also says "Positional and keyword arguments are 
supported."
   
   The reverse case: a callable instance whose method is `def __call__(self)` 
passes the check, because `getfullargspec` reports `args=['self']`, and then 
fails at runtime.
   
   `pandas_udf` / `arrow_udf` have the same `**kwargs` limitation, so this is a 
nit. Adding `and argspec.varkw is None` would make the check consistent with 
the documented keyword support.



##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/PythonUDF.scala:
##########
@@ -51,7 +51,8 @@ object PythonUDF {
     PythonEvalType.SQL_SCALAR_PANDAS_UDF,
     PythonEvalType.SQL_SCALAR_PANDAS_ITER_UDF,
     PythonEvalType.SQL_SCALAR_ARROW_UDF,
-    PythonEvalType.SQL_SCALAR_ARROW_ITER_UDF
+    PythonEvalType.SQL_SCALAR_ARROW_ITER_UDF,
+    PythonEvalType.SQL_SCALAR_ARROW_INPROCESS_UDF

Review Comment:
   **[Low, performance] Adding 258 to `SCALAR_TYPES` makes `CollapseProject` 
force-inline producers that contain in-process UDFs.**
   
   `CollapseProject.mergeProjectExpressions` (currently `Optimizer.scala` 
~L1794-1805) classifies an alias as `mustInlines` when its child contains a 
scalar Python UDF whose eval type appears in the upper project (`alwaysInline 
|| containsUDF`). It does this to enable UDF chaining, and it skips the usual 
reference-count and cheapness checks. But 
`ExtractPythonUDFs.shouldExtractUDFExpressionTree` never chains 258, so for 
in-process UDFs this forced inline only duplicates work.
   
   Example: `df.select((inproc1("x") + F.length(F.sha2("s", 
512))).alias("a")).select(inproc2("a"), F.col("a"))`. Without 258 in 
`SCALAR_TYPES`, `a` is referenced twice and not cheap, so it would never be 
inlined and `sha2` would run once per row. Now it is force-inlined, and `p0 + 
length(sha2(s))` is computed twice per row: once in the EvalPython input 
projection and once in the upper Project.
   
   Suggestion: exclude 258 from the chaining-driven inline decision.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalPythonEvaluatorFactory.scala:
##########
@@ -0,0 +1,233 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.util.UUID
+
+import scala.collection.mutable.ArrayBuffer
+import scala.jdk.CollectionConverters._
+
+import org.apache.arrow.c.{ArrowArray, ArrowSchema}
+import org.apache.arrow.util.AutoCloseables
+import org.apache.arrow.vector.VectorSchemaRoot
+
+import org.apache.spark.TaskContext
+import org.apache.spark.api.python.ChainedPythonFunctions
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{Attribute, PythonUDF}
+import org.apache.spark.sql.execution.arrow.ArrowWriter
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.execution.python.EvalPythonExec.ArgumentMetadata
+import org.apache.spark.sql.types.StructType
+import org.apache.spark.sql.util.ArrowUtils
+import org.apache.spark.sql.vectorized.{ArrowColumnVector, ColumnarBatch, 
ColumnVector}
+import org.apache.spark.util.Utils
+
+/**
+ * Evaluates scalar Python UDFs using Arrow CDI in the executor process. Only 
UDF arguments
+ * are converted to Arrow. Original rows are buffered in a spillable queue and 
joined with
+ * the results. Each batch owns its Arrow buffers so Python can safely retain 
input arrays.
+ */
+class InProcessArrowEvalPythonEvaluatorFactory(
+    childOutput: Seq[Attribute],
+    udfs: Seq[PythonUDF],
+    output: Seq[Attribute],
+    batchSize: Int,
+    maxBytes: Long,
+    timeZoneId: String,
+    largeVarTypes: Boolean,
+    hideTraceback: Boolean,
+    simplifiedTraceback: Boolean,
+    tracebackWithLocals: Boolean,
+    metrics: Map[String, SQLMetric])
+  extends EvalPythonEvaluatorFactory(childOutput, udfs, output) {
+
+  override protected def evaluate(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputSchema: StructType,
+      context: TaskContext): Iterator[InternalRow] = {
+    ArrowUtils.failDuplicatedFieldNames(inputSchema)
+    val functions = funcs.map { case (chain, _) =>
+      require(chain.funcs.size == 1, "In-process UDF chains must use separate 
evaluation nodes")
+      chain.funcs.head
+    }
+    val inputOrdinals = argMetas.map(_.map(_.offset))
+    def checkCancellation(): Unit = context.killTaskIfInterrupted()
+
+    val expectedFields = udfs.map { udf =>
+      ArrowUtils.toArrowField("result", udf.dataType, true, timeZoneId, 
largeVarTypes)
+    }
+    val processingTime = new 
InProcessArrowEvalPythonEvaluatorFactory.NanosecondTimer(
+      metrics("pythonProcessingTime"))
+    val initTime = new 
InProcessArrowEvalPythonEvaluatorFactory.NanosecondTimer(
+      metrics("pythonInitTime"))
+    val arrowSchema = ArrowUtils.toArrowSchema(inputSchema, timeZoneId, 
largeVarTypes)
+    var runtime: InProcessPythonRuntime.InterpreterSession = null
+    val handles = functions.map(_ => UUID.randomUUID().toString)
+    var registered = false
+    var writer: ArrowWriter = null
+    val results = ArrayBuffer.empty[ArrowColumnVector]
+    var closed = false
+    val startedAt = System.nanoTime()
+
+    def closeBatch(): Unit = {
+      val resources = ArrayBuffer.empty[AutoCloseable]
+      resources ++= results
+      results.clear()
+      if (writer != null) {
+        resources += writer.root
+        writer = null
+      }
+      AutoCloseables.close(resources.asJava)
+    }
+
+    def close(): Unit = {
+      if (!closed) {
+        closed = true
+        metrics("pythonTotalTime") += (System.nanoTime() - startedAt) / 1000000
+        Utils.tryWithSafeFinally {
+          closeBatch()
+        } {
+          if (registered) runtime.release(handles)
+        }
+      }
+    }
+
+    context.addTaskCompletionListener[Unit](_ => close())
+
+    new Iterator[InternalRow] {
+      private var batchIter: Iterator[InternalRow] = Iterator.empty
+
+      override def hasNext: Boolean = {
+        checkCancellation()
+        val available = !closed && (batchIter.hasNext || rows.hasNext)
+        if (!available) close()
+        available
+      }
+
+      override def next(): InternalRow = {
+        if (!hasNext) throw new NoSuchElementException("End of in-process UDF 
input")
+        try {
+          if (!batchIter.hasNext) {
+            closeBatch()
+            if (!registered) {
+              runtime = InProcessPythonRuntime.currentSession
+              // Mark before registering so failure after any registration 
still cleans up.
+              registered = true
+              val start = System.nanoTime()

Review Comment:
   **[Low] `pythonInitTime` includes time spent queued behind other tasks.**
   
   This timer runs on the task thread around `runtime.register`, and `register` 
blocks in `onInterpreterThread` until the single interpreter thread picks the 
call up. `invoke`, by contrast, starts its timer inside the interpreter thread. 
So with multiple tasks per executor (supported and tested, e.g. `local[2]` with 
`spark.task.cpus=0.5`), a task that starts while another task is running a long 
UDF batch reports that other task's Python time as its own init time.
   
   Relatedly, `pythonTotalTime` is charged from the completion listener even 
when the iterator is never consumed. The worker path records it only from 
end-of-stream timing data.
   
   Suggestion: measure the registration time inside the interpreter thread, as 
`invoke` already does.



##########
python/pyspark/inprocess/udf.py:
##########
@@ -0,0 +1,225 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+"""
+Python API for in-process UDF registration.
+
+Usage::
+
+    import pyarrow.compute as pc
+    from pyspark.inprocess import inprocess_udf
+    from pyspark.sql.types import LongType
+
+    @inprocess_udf(return_type=LongType())
+    def double(x):
+        # x is a pa.Array; return a pa.Array
+        return pc.multiply(x, 2)
+
+    df.select(double(df.value)).show()
+"""
+
+import io
+import sys
+from functools import update_wrapper
+from inspect import getfullargspec
+from typing import Any, Callable, Optional, Union
+
+from pyspark import Accumulator, Broadcast, cloudpickle
+from pyspark.errors import PySparkTypeError, PySparkValueError
+from pyspark.sql.column import Column
+from pyspark.sql.types import DataType, _parse_datatype_string
+from pyspark.util import PythonEvalType
+
+
+class _InProcessPickler(cloudpickle.CloudPickler):
+    def reducer_override(self, obj: Any) -> Any:
+        if isinstance(obj, (Broadcast, Accumulator)):
+            raise TypeError("In-process UDFs do not support Spark broadcasts 
or accumulators")
+        return super().reducer_override(obj)
+
+
+def _serialize_udf(func: Callable, return_type: DataType) -> bytes:
+    buffer = io.BytesIO()
+    _InProcessPickler(buffer).dump((func, return_type))
+    return buffer.getvalue()
+
+
+class InProcessUDFWrapper:
+    """
+    Wraps a Python function as an in-process UDF.
+
+    Returned by ``@inprocess_udf``. Calling an instance with Spark ``Column``
+    arguments creates a ``Column`` expression backed by ``PythonUDF``
+    on the JVM side.
+    """
+
+    def __init__(
+        self, func: Callable, return_type: Union[DataType, str], 
deterministic: bool = True
+    ) -> None:
+        if not isinstance(return_type, (DataType, str)):
+            raise PySparkTypeError(
+                errorClass="NOT_EXPECTED_TYPE",
+                messageParameters={
+                    "expected_type": "DataType or str",
+                    "arg_name": "return_type",
+                    "arg_type": type(return_type).__name__,
+                },
+            )
+        self._return_type = return_type
+        self._parsed_return_type: Optional[DataType] = None
+        self.evalType = PythonEvalType.SQL_SCALAR_ARROW_INPROCESS_UDF
+        self._deterministic: bool = deterministic
+        self._name: str = getattr(func, "__name__", "inprocess_udf")
+
+        argspec = getfullargspec(func)
+        if not argspec.args and argspec.varargs is None and not 
argspec.kwonlyargs:
+            raise PySparkValueError(
+                errorClass="INVALID_PANDAS_UDF",
+                messageParameters={"detail": "0-arg inprocess_udfs are not 
supported."},
+            )
+        self._func = func
+        self._serialized: Optional[bytes] = None
+        update_wrapper(self, func, updated=())
+
+    @property
+    def func(self) -> Callable:
+        return self._func
+
+    @property
+    def returnType(self) -> DataType:
+        if self._parsed_return_type is None:
+            parsed = (
+                _parse_datatype_string(self._return_type)
+                if isinstance(self._return_type, str)
+                else self._return_type
+            )
+            from pyspark.sql.udf import UserDefinedFunction
+
+            UserDefinedFunction._check_return_type(parsed, 
PythonEvalType.SQL_SCALAR_ARROW_UDF)

Review Comment:
   **[Medium] Struct return types with duplicate field names pass the driver 
check and fail on every task.**
   
   `_check_return_type(..., SQL_SCALAR_ARROW_UDF)` calls 
`to_arrow_type(returnType, timezone="UTC")`, and that defaults to 
`error_on_duplicated_field_names_in_struct=False`. `_inprocess_register` on the 
executor passes `True`, and the JVM side (`ArrowUtils.toArrowField` for 
`expectedFields`) does no such check.
   
   So `inprocess_udf(StructType([StructField("a", LongType()), StructField("a", 
LongType())]))(f)(col)` builds and plans normally. Every non-empty partition 
then fails with `DUPLICATED_FIELD_NAME_IN_ARROW_STRUCT` (as a 
`PythonException`), retried up to `spark.task.maxFailures`. This contradicts 
the guide ("Unsupported return types are rejected on the driver before the 
function is serialized or tasks start").
   
   Suggestion: here, also call `to_arrow_type(parsed, timezone="UTC", 
error_on_duplicated_field_names_in_struct=True)`, and add a test for a return 
struct with duplicate names.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,322 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+
+"""Arrow CDI entry points called on the executor's dedicated JEP interpreter 
thread.
+
+Functions are registered once per task and released when that task finishes. 
Calls
+pass only a handle and CDI addresses, so large closures are not copied per 
batch.
+"""
+
+import sys
+from typing import Any, Callable, Iterable, Optional, Sequence
+
+import pyarrow as pa
+import pyarrow.compute as pc
+
+from pyspark import cloudpickle
+from pyspark.errors import PySparkRuntimeError
+from pyspark.sql.pandas.types import to_arrow_type
+from pyspark.util import _format_exception
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+NullChecker = Callable[[pa.Array], None]
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType, NullChecker, 
bool, bool, bool]] = {}
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    timezone: str,
+    python_version: str,
+    large_var_types: bool = False,
+    hide_traceback: bool = False,
+    simplified_traceback: bool = False,
+    traceback_with_locals: bool = False,
+) -> None:
+    try:
+        embedded_version = "%d.%d" % sys.version_info[:2]
+        if python_version != embedded_version:
+            raise PySparkRuntimeError(
+                errorClass="PYTHON_VERSION_MISMATCH",
+                messageParameters={
+                    "worker_version": embedded_version,
+                    "driver_version": python_version,
+                },
+            )
+        # JEP exposes direct ByteBuffers through the buffer protocol. Unpickle 
a separate
+        # function per task without iterating over a PyJArray one JNI call per 
byte.
+        # Carry the type with the closure so driver-defined UDTs need no 
module import.
+        func, return_type = cloudpickle.loads(memoryview(serialized_udf))
+        expected_type = to_arrow_type(
+            return_type,
+            timezone=timezone,
+            prefers_large_types=large_var_types,
+            error_on_duplicated_field_names_in_struct=True,
+        )
+        if large_var_types:
+            expected_type = _large_binary_type(expected_type)

Review Comment:
   **[Design] The expected result Arrow type is derived twice, once in Python 
and once in the JVM.**
   
   The executor side uses `to_arrow_type` on a return type pickled together 
with the closure. The JVM uses `ArrowUtils.toArrowField` for `expectedFields`. 
This PR already patches the resulting divergences on both sides:
   - `_large_binary_type` for Variant/Geometry/Geography binary children;
   - JSON metadata whitespace and TIME precision tags, handled by `sameLayout` 
ignoring metadata plus rebuilding the vector from the declared field;
   - `_with_schema` renaming.
   
   Each patch has its own test. Any future layout change on one side that isn't 
mirrored on the other would make in-process UDFs fail at runtime 
(`require(sameLayout(...))` / `TypeError`), while worker Arrow UDFs keep 
working because they compare Spark types.
   
   A more robust option is to make the JVM the single source of truth: export 
`expectedFields(i)` once at registration with `Data.exportField`, and import it 
in Python with `pa.Field._import_from_c` / `pa.DataType._import_from_c`. That 
would also remove the pickled return type, the `timezone` / `large_var_types` 
register arguments, and the executor-side `to_arrow_type` / 
`_large_binary_type`.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/ArrowEvalPythonExec.scala:
##########
@@ -66,6 +66,7 @@ private[spark] class BatchIterator[T](iter: Iterator[T], 
batchSize: Int)
  * Following eval types are supported:
  *
  * <ul>
+ *   <li> SQL_SCALAR_ARROW_INPROCESS_UDF for an embedded scalar Arrow UDF

Review Comment:
   **[Nit] This bullet says `ArrowEvalPythonExec` supports 
`SQL_SCALAR_ARROW_INPROCESS_UDF`, but it rejects it.**
   
   `supportedPythonEvalTypes` (L169-180) doesn't include 258, so the 
constructor throws `SparkException.internalError("Unexpected eval type 258")`. 
In-process UDFs are routed correctly only because the new first case in 
`SparkStrategies.PythonEvals` sends them to `InProcessArrowEvalPythonExec`. 
Please drop this bullet.



-- 
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