viirya commented on code in PR #58978:
URL: https://github.com/apache/spark/pull/58978#discussion_r4112403863


##########
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:
   Changed initialization to `PyConfig.isolated().setUseEnvironment(false)`, 
while retaining the fixed hash seed. A fresh-JVM test sets both 
`PYTHONFAULTHANDLER` and `PYTHONDEVMODE` and verifies isolated mode, ignored 
environment settings, and disabled faulthandler. Required Python paths are 
restored explicitly.



##########
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:
   The embedded runtime now merges `PythonUtils.sparkPythonPath` with the 
executor process's `PYTHONPATH` and places both ahead of configured 
site-packages. This also restores the paths that isolated initialization no 
longer picks up automatically.
   
   Added a fresh-JVM test without `SPARK_HOME`, using Spark archives through 
`PYTHONPATH` and a conflicting PySpark package in configured site-packages. The 
guide now explicitly documents that the entire process `PYTHONPATH` takes 
precedence. Existing worker startup is unchanged.



##########
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:
   Updated the Kubernetes example to use 
`spark.executor.extraJavaOptions=-Djava.library.path=...` and explained why it 
is needed for images whose entrypoint does not propagate 
`spark.executor.extraLibraryPath`. This was checked against the launcher code; 
I haven't validated it on a Kubernetes cluster.



##########
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:
   Split main-interpreter setup from shared-interpreter configuration tracking. 
Shared configuration is marked successful only after JEP's 
`configureInterpreter` completes, so a failed initial configuration can be 
replaced without repeating `setInitParams`.
   
   A small `SharedInterpreter` subclass also closes the native handle if 
configuration fails during construction. A fresh-JVM regression test verifies 
that initialization succeeds after correcting a site-packages path that 
initially cannot provide JEP.



##########
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:
   Added ASCII escaping before exception text crosses JNI, with explicit NUL 
escaping because `backslashreplace` alone leaves NUL unchanged. Bootstrap 
errors now use `ascii(error)` as well. Pure-Python and JEP integration tests 
cover non-BMP characters, lone surrogates, and embedded NUL. These characters 
remain visible as escape sequences in the reported error.



##########
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:
   Added driver-side rejection of `spark.executor.pyspark.memory` and 
`spark.sql.pyspark.udf.profiler`, alongside `spark.pythonWorkerEnv.*`. The 
guide explains the limitations and points users to executor memory settings or 
worker-based UDFs. Tests cover both additional settings.



##########
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:
   Plugin initialization now performs a small CDI field export/import round 
trip, checking both the JAR and native library before accepting tasks. CDI 
class resolution happens inside the guarded call, and the initialization error 
names the required dependencies. A fresh-JVM test without the CDI JAR verifies 
the early failure and installation hint.



##########
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:
   Added the strict driver-side `to_arrow_type` validation with 
`timezone="UTC"` and `error_on_duplicated_field_names_in_struct=True`. Tests 
cover duplicate names in both top-level and nested structs and verify rejection 
before the serialized command is cached.



##########
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:
   Added `INVALID_SPARK_CONFIG.UNSUPPORTED_IN_PROCESS_PYTHON_UDF` and moved 
execution-time configuration validation into the dedicated exec's `doExecute`, 
before task submission. Tests cover configuration changes after Column 
construction with partition evaluators both enabled and disabled.
   
   The internal invariant checks in the bridge and evaluator now use 
`SparkException.internalError`.



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