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


##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/PythonUDF.scala:
##########
@@ -48,7 +48,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:
   Now that 258 is a scalar eval type, a Spark Connect client can send a 
`python_udf` with `eval_type = 258`, because `SparkConnectPlanner` does not 
validate eval types, and it gets planned into 
`InProcessArrowEvalPythonEvaluatorFactory`. The Connect command is the pickled 
`(func, returnType)` tuple. With the plugin enabled, it is unpickled in the 
executor's shared interpreter and then fails with `'tuple' object is not 
callable`. This path also skips per-session isolation (pythonIncludes, envVars, 
job artifact UUID). Before this PR, such a request failed with an internal 
error and no Python code ran. Should Connect reject 258 explicitly?
   
   Relatedly, `register` in `python/pyspark/sql/connect/udf.py` has no 
`InProcessUDFWrapper` check. On Connect, `spark.udf.register("f", 
inprocess_fn)` succeeds as a StringType batched UDF, and every later call fails 
with "No active SparkContext".



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowBridge.scala:
##########
@@ -0,0 +1,116 @@
+/*
+ * 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.arrow.c.{ArrowArray, ArrowSchema, Data}
+import org.apache.arrow.memory.util.MemoryUtil
+import org.apache.arrow.vector.FieldVector
+import org.apache.arrow.vector.types.pojo.Field
+
+import org.apache.spark.sql.util.ArrowUtils
+import org.apache.spark.sql.vectorized.ArrowColumnVector
+import org.apache.spark.util.Utils
+
+/**
+ * Bridges JVM Arrow column buffers with Python PyArrow arrays for in-process 
UDF execution.
+ *
+ * Both input and output paths use the Arrow C Data Interface (CDI) for 
zero-copy transfer.
+ *
+ * Input path (JVM to Python, zero-copy via CDI):
+ *   JVM pre-allocates [[ArrowArray]] and [[ArrowSchema]] C structs and 
exports each input
+ *   [[FieldVector]] into them via [[Data.exportVector]]. The native addresses 
are passed to
+ *   Python. Python calls ``pa.Array._import_from_c(array_ptr, schema_ptr)`` 
to wrap the
+ *   same Arrow buffers as a PyArrow array -- no memcpy. When Python GCs the 
array, the CDI
+ *   release callback decrements the buffer reference counts; the JVM 
[[FieldVector]] retains
+ *   its own reference. Each batch uses new vectors; closing the old vectors 
releases only the
+ *   JVM's references, leaving any arrays retained by Python valid and 
unchanged.
+ *
+ * Output path (Python to JVM, zero-copy via CDI):
+ *   JVM pre-allocates [[ArrowArray]] and [[ArrowSchema]] C structs. Python 
calls
+ *   ``arr._export_to_c(array_ptr, schema_ptr)`` to fill those structs 
in-place. The JVM
+ *   calls [[Data.importIntoVector]] to reconstruct the [[FieldVector]] 
without copying. When the
+ *   imported [[FieldVector]] is closed, Arrow Java invokes PyArrow's CDI 
release callback,
+ *   decrementing the Python array refcount and allowing garbage collection.
+ *
+ * The runtime validates the returned schema before ArrowColumnVector reads 
the buffers.
+ */
+private[python] object InProcessArrowBridge {
+
+  /**
+   * Export a [[FieldVector]] to pre-allocated Arrow C Data Interface structs.
+   *
+   * Fills ``outArray`` and ``outSchema`` with the CDI representation of 
``vector``.
+   * The export is zero-copy: ``outArray``'s buffer pointers reference the 
same off-heap
+   * memory as ``vector``. The CDI release callback (invoked when the 
Python-side imported
+   * array is GC'd) decrements the buffer reference counts; the 
[[FieldVector]] continues
+   * to hold its own reference.
+   *
+   * Caller must release any unconsumed exports and close both structs on 
every exit path.
+   */
+  def exportColumn(vector: FieldVector, outArray: ArrowArray, outSchema: 
ArrowSchema): Unit =
+    Data.exportVector(ArrowUtils.rootAllocator, vector, null, outArray, 
outSchema)
+
+  /**
+   * Reconstruct an [[ArrowColumnVector]] from JVM-allocated Arrow C Data 
Interface structs.
+   *
+   * The JVM pre-allocates [[ArrowArray]] and [[ArrowSchema]] before invoking 
Python.
+   * Python fills them via ``arr._export_to_c(array_ptr, schema_ptr)``. This 
method
+   * calls [[Data.importIntoVector]] to wrap Python's Arrow buffers 
(zero-copy).
+   *
+   * Lifecycle:
+   *  - [[Data.importIntoVector]] internally calls 
``ArrayImporter.importArray()``, which
+   *    moves the struct snapshot through a non-owning wrapper, leaving the 
caller's struct
+   *    storage alive for cleanup, and wraps the data buffers via
+   *    ``ReferenceCountedArrowArray`` (ForeignAllocation, zero-copy).
+   *  - Data.importField releases and closes a non-owning schema wrapper too.
+   *    The caller closes the original struct storage.
+   *  - When the returned [[ArrowColumnVector]] is closed, the reference count 
drops to
+   *    zero, PyArrow's C ``release`` callback is invoked, and the Python 
array is GC'd.
+   */
+  private def checkOffsets(array: ArrowArray): Unit = {
+    val snapshot = array.snapshot()
+    require(snapshot.offset == 0L, "In-process UDF returned an unsupported 
Arrow CDI offset")
+    (0L until snapshot.n_children).foreach { i =>
+      checkOffsets(ArrowArray.wrap(MemoryUtil.getLong(snapshot.children + i * 
8L)))
+    }
+    if (snapshot.dictionary != 0L) 
checkOffsets(ArrowArray.wrap(snapshot.dictionary))
+  }
+
+  def cdiToColumn(
+      arrowArray: ArrowArray,
+      arrowSchema: ArrowSchema,
+      expected: Option[Field] = None): ArrowColumnVector = {
+    checkOffsets(arrowArray)
+    val field = Data.importField(
+      ArrowUtils.rootAllocator, ArrowSchema.wrap(arrowSchema.memoryAddress()), 
null)
+    expected.foreach { declared =>
+      require(field.getType == declared.getType && field.getChildren == 
declared.getChildren &&

Review Comment:
   This exact `Field` comparison (children and metadata included) rejects valid 
results where the JVM and Python encodings differ:
   - Struct fields with Spark metadata: the JVM side 
(`ArrowUtils.toArrowMetaData`) uses `Metadata.json`, which is compact 
(`{"comment":"c"}`), while Python's `to_arrow_metadata` uses `json.dumps` 
(`{"comment": "c"}`). `_validate_result` passes because pyarrow type equality 
ignores metadata, so even an identity UDF declared with 
`StructType([StructField("a", LongType(), True, {"comment": "c"})])` fails 
every batch here. This is easy to hit when reusing a table schema that has 
column comments.
   - TimeType and the nanos timestamp types (see the comment in the evaluator 
factory).
   
   Could we compare Spark types instead (e.g. `fromArrowField` + 
`equalsIgnoreCompatibleCollation`, as `ArrowEvalPythonEvaluatorFactory` does), 
or at least compare parsed metadata instead of raw strings?



##########
python/benchmarks/bench_inprocess_udf.py:
##########
@@ -0,0 +1,137 @@
+#
+# 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.
+#
+
+"""End-to-end in-process, worker Arrow, and pandas UDF benchmarks.
+
+See README.md for the required Spark build and JEP launch environment. These
+measure steady-state queries, including JVM row/Arrow conversion and Python
+execution. Worker Arrow UDFs are the primary baseline and use the same Arrow
+operations as in-process UDFs. The supplementary pandas baseline also includes
+pandas conversion costs; neither comparison isolates IPC overhead alone.
+Historical standalone-script timings are a separate baseline.
+"""
+
+from importlib.util import find_spec
+
+
+class InProcessUDFTimeBench:
+    # One query per sample, with explicit full-query warmup in setup.
+    number = 1
+    rounds = 1
+    repeat = 5
+    warmup_time = 0
+    timeout = 300
+    params = [
+        ["arrow", "inprocess", "pandas"],
+        [
+            ("narrow", 100_000),
+            ("narrow", 1_000_000),
+            ("narrow", 5_000_000),
+            ("wide", 1_000_000),
+            ("wide", 5_000_000),
+            ("wide", 10_000_000),
+            ("short_string", 1_000_000),
+            ("short_string", 5_000_000),
+            ("short_string", 10_000_000),
+            ("long_string", 500_000),
+            ("long_string", 1_000_000),
+            ("long_string", 2_000_000),
+        ],
+    ]
+    param_names = ["udf_type", "workload"]
+
+    def setup(self, udf_type, workload):
+        # JEP cannot be imported from standalone CPython. Check availability
+        # without loading it; broken native/JVM setup must fail, not be 
skipped.
+        if udf_type == "inprocess" and find_spec("jep") is None:
+            raise NotImplementedError("Install JEP and configure its JVM 
launch paths")
+
+        import pyarrow.compute as pc
+        from pyspark.sql import SparkSession
+        from pyspark.sql.functions import arrow_udf, col, lpad, pandas_udf
+        from pyspark.sql.types import LongType, StringType
+
+        use_arrow = udf_type != "pandas"
+        scenario, n_rows = workload
+        n_cols = 10 if scenario == "wide" else 1
+        batch_size = {"narrow": 10_000, "wide": 1_000_000}.get(scenario, 
100_000)
+        self.spark = (
+            SparkSession.builder.master("local[1]")

Review Comment:
   This session never sets 
`spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin`, 
and neither does `PYSPARK_SUBMIT_ARGS` in the README. The lazy interpreter 
initialization has been removed, so the `inprocess` cases now fail in the setup 
warmup with `In-process Python is not running; initialize the executor plugin 
first`.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalPythonEvaluatorFactory.scala:
##########
@@ -0,0 +1,212 @@
+/*
+ * 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,
+    metrics: Map[String, SQLMetric])
+  extends EvalPythonEvaluatorFactory(childOutput, udfs, output) {
+
+  private val returnTypes = udfs.map(_.dataType.json)
+
+  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 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()
+              functions.indices.foreach { i =>
+                val func = functions(i)
+                runtime.register(handles(i), func.command.toArray, 
returnTypes(i),
+                  timeZoneId, func.pythonVer, largeVarTypes)
+              }
+              metrics("pythonInitTime") += (System.nanoTime() - start) / 
1000000
+            }
+            val root = VectorSchemaRoot.create(arrowSchema, 
ArrowUtils.rootAllocator)
+            writer = try {
+              ArrowWriter.create(root)
+            } catch {
+              case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
root.close() }
+            }
+            var count = 0
+            while (rows.hasNext && (batchSize <= 0 || count < batchSize) &&
+                (count == 0 || maxBytes <= 0 || writer.sizeInBytes() < 
maxBytes)) {
+              checkCancellation()
+              writer.write(rows.next())
+              count += 1
+            }
+            writer.finish()
+            metrics("pythonDataSent") += writer.sizeInBytes()
+
+            handles.indices.foreach { udfIndex =>
+              val handle = handles(udfIndex)
+              val ordinals = inputOrdinals(udfIndex)
+              checkCancellation()
+              // Register each acquired resource immediately, including 
partially exported
+              // inputs and results of earlier UDFs if a later UDF throws.
+              val structs = ArrayBuffer.empty[AutoCloseable]
+              def array(): ArrowArray = {
+                val value = ArrowArray.allocateNew(ArrowUtils.rootAllocator)
+                structs += new AutoCloseable {
+                  override def close(): Unit =
+                    Utils.tryWithSafeFinally {
+                      if (value.snapshot().release != 0L) value.release()
+                    } { value.close() }
+                }
+                value
+              }
+              def schema(): ArrowSchema = {
+                val value = ArrowSchema.allocateNew(ArrowUtils.rootAllocator)
+                structs += new AutoCloseable {
+                  override def close(): Unit =
+                    Utils.tryWithSafeFinally {
+                      if (value.snapshot().release != 0L) value.release()
+                    } { value.close() }
+                }
+                value
+              }
+              Utils.tryWithSafeFinally {
+                val inArrays = ordinals.map(_ => array())
+                val inSchemas = ordinals.map(_ => schema())
+                val outArray = array()
+                val outSchema = schema()
+                ordinals.indices.foreach { i =>
+                  InProcessArrowBridge.exportColumn(
+                    writer.root.getVector(ordinals(i)), inArrays(i), 
inSchemas(i))
+                }
+                metrics("pythonProcessingTime") += runtime.invoke(
+                  handle,
+                  inArrays.map(_.memoryAddress()).toArray,
+                  inSchemas.map(_.memoryAddress()).toArray,
+                  outArray.memoryAddress(), outSchema.memoryAddress(),
+                  count, argMetas(udfIndex).map(_.name.getOrElse("")))
+                val expected = ArrowUtils.toArrowField(

Review Comment:
   For `TimeType` and the nanos timestamp types, `toArrowField` adds precision 
metadata (`SPARK::time::precision` / `SPARK::timestampNanos::precision`) to 
this top-level field. `pa.Array._export_to_c` exports a type-only schema with 
no field metadata. These results pass the Python check and then always fail the 
`require` in `cdiToColumn` (`{}` vs `{SPARK::time::precision=6}`). Nested 
`TimeType` fails too, because Python does not tag it.
   
   Minor: `expected` does not change between batches, so it could be computed 
once per UDF next to `returnTypes`.



##########
python/pyspark/inprocess/bridge.py:
##########
@@ -0,0 +1,38 @@
+#
+# 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.
+#
+
+"""

Review Comment:
   This module contains only a docstring and is not imported anywhere. It 
duplicates the `InProcessArrowBridge` scaladoc and has already drifted: it says 
`Data.importVector`, while the code calls `Data.importIntoVector`. Can we 
remove it?



##########
docs/sql-pyspark-inprocess-udf.md:
##########
@@ -0,0 +1,596 @@
+---
+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. 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.
+
+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. 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. `ArrowEvalPythonExec` 
selects
+an in-process evaluator factory for this evaluation type, reusing the 
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. Extra site-packages paths are
+processed with `site.addsitedir` before loading the runtime bridge, including 
`.pth`
+files. Configured directories and newly discovered `.pth` paths precede 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`. SQL registration through
+`spark.udf.register` is not supported and is rejected at registration time.
+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. 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 evaluated once at UDF definition time and shipped with the function 
to every executor:
+
+```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 are already on `sys.path`, so no extra configuration is needed.
+
+```bash
+python3 -m venv .venv
+.venv/bin/pip install "jep>=4.3.2" pyarrow cloudpickle pyspark
+source .venv/bin/activate
+```
+
+You must also make the jep native library discoverable by the JVM:
+
+```bash
+# macOS
+export DYLD_LIBRARY_PATH="$(python3 -c 'import jep; import os; 
print(os.path.dirname(jep.__file__))')"
+
+# Linux
+export LD_LIBRARY_PATH="$(python3 -c 'import jep; import os; 
print(os.path.dirname(jep.__file__))')"
+```
+
+### 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 at task
+launch time. 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

Review Comment:
   The YARN and Kubernetes (Option A/B) examples do not put the JEP JAR or the 
`arrow-c-data` JAR on the driver/executor classpath. `InProcessPythonRuntime` 
and `InProcessArrowBridge` are loaded by the system classloader, so `--jars` 
would not help either. `spark.{driver,executor}.extraClassPath`, or copying the 
JARs into `$SPARK_HOME/jars`, seems required. Otherwise plugin init fails with 
`NoClassDefFoundError: jep/...`.
   
   Also, lines 283 and 286 use `python3 -c 'import jep'`. With JEP 4.3.x this 
raises `ImportError` in standalone Python, so the exported library path becomes 
empty. The `find_spec("jep").origin` approach from the Dockerfile example would 
work instead.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, TimeoutException, 
TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepException, SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.PythonException
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.{ThreadUtils, 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 = _
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      require(active.isRunning && active.sitePackages == sitePackages,
+        "In-process Python is stopping or already initialized with different 
sitePackages")
+    } else {
+      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) {
+    private val executor = 
ThreadUtils.newDaemonSingleThreadExecutor("inprocess-python")
+    @volatile private var running = true
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+
+    def isRunning: Boolean = running
+    def isTerminated: Boolean = executor.isTerminated
+
+    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()

Review Comment:
   Regular Python workers always run with `PYTHONHASHSEED` (default `"0"`, set 
in `sc.environment` by `context.py` and passed through the function's envVars). 
Here the interpreter is created with JEP's default `PyConfig`, and 
`InProcessPythonUDFBuilder` passes empty envVars. On YARN, K8s, and standalone, 
the executor JVM environment has no `PYTHONHASHSEED`, so the hash seed is 
random per executor process. A `deterministic=True` UDF that uses `hash(s) % n` 
or set iteration order can then return different results when a task is retried 
or recomputed. Could we set a fixed hash seed at initialization (e.g. via 
`MainInterpreter.setInitParams` with a `PyConfig` hash seed), or at least 
document `spark.executorEnv.PYTHONHASHSEED`?



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,196 @@
+#
+# 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
+import traceback as _traceback
+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.sql.types import _parse_datatype_json_string
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType]] = {}
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    return_type_json: str,
+    timezone: str,
+    python_version: str,
+    large_var_types: 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.
+        func = cloudpickle.loads(memoryview(serialized_udf))
+        expected_type = to_arrow_type(

Review Comment:
   With `spark.sql.execution.arrow.useLargeVarTypes=true`, `to_arrow_type` 
still hard-codes `pa.binary()` children for Variant, Geometry, and Geography. 
The JVM expected field, and the exported inputs, use `LargeBinary`. So no 
result can pass both this check and the JVM `require` in `cdiToColumn`. For 
example, an identity UDF on a Variant column fails here with a TypeError. The 
worker path is not affected because it has no strict type equality check.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,196 @@
+#
+# 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
+import traceback as _traceback
+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.sql.types import _parse_datatype_json_string
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType]] = {}
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    return_type_json: str,
+    timezone: str,
+    python_version: str,
+    large_var_types: 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.
+        func = cloudpickle.loads(memoryview(serialized_udf))
+        expected_type = to_arrow_type(
+            _parse_datatype_json_string(return_type_json),
+            timezone=timezone,
+            prefers_large_types=large_var_types,
+            error_on_duplicated_field_names_in_struct=True,
+        )
+        _udfs[handle] = (func, expected_type)
+    except BaseException:
+        # In JEP, an uncaught SystemExit can terminate the entire executor JVM.
+        raise RuntimeError(_UDF_TRACEBACK_SENTINEL + _traceback.format_exc()) 
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 _check_nested_nulls(array: pa.Array, expected_type: pa.DataType) -> None:
+    def check_field(values: pa.Array, field: pa.Field) -> None:
+        if not field.nullable and values.null_count:
+            raise ValueError(f"In-process UDF returned nulls in non-nullable 
field {field.name}")
+        _check_nested_nulls(values, field.type)
+
+    if pa.types.is_struct(expected_type):
+        # Children under a null parent do not contribute values to the result.
+        visible = pc.filter(array, pc.is_valid(array))
+        for i, field in enumerate(expected_type):
+            check_field(visible.field(i), field)
+    elif pa.types.is_list(expected_type) or 
pa.types.is_large_list(expected_type):
+        check_field(pc.list_flatten(array), expected_type.value_field)
+    elif pa.types.is_map(expected_type):
+        visible = pa.concat_arrays([pc.filter(array, pc.is_valid(array))])
+        check_field(visible.keys, expected_type.key_field)
+        check_field(visible.items, expected_type.item_field)
+
+
+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) -> 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()
+    _check_nested_nulls(result, expected_type)
+    # 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.
+    """
+    try:
+        udf_func, expected_type = _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)
+        result._export_to_c(int(output_array_ptr), int(output_schema_ptr))
+    except BaseException:
+        raise RuntimeError(_UDF_TRACEBACK_SENTINEL + _traceback.format_exc()) 
from None

Review Comment:
   This always uses `traceback.format_exc()`, so 
`spark.sql.execution.pyspark.udf.hideTraceback.enabled` and 
`simplifiedTraceback.enabled` are ignored for in-process UDFs. The worker path 
honors both through `SPARK_HIDE_TRACEBACK` / `SPARK_SIMPLIFIED_TRACEBACK` and 
`handle_worker_exception`. The full traceback also ends up in the 
`JepException` cause on the JVM side. The same applies at line 71.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, TimeoutException, 
TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepException, SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.PythonException
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.{ThreadUtils, 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 = _
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      require(active.isRunning && active.sitePackages == sitePackages,

Review Comment:
   `shutdown()` returns after the 5s wait, but `active` still points at the old 
session. If a call is still running (not necessarily hung), creating a new 
SparkContext in the same JVM fails this `require`. This happens in local mode 
or in a notebook that calls `spark.stop()` and then `getOrCreate()`. The plugin 
then logs "Verify that: (1) libjep.so ... (2) jep.jar ...", which sends users 
to check their installation instead of pointing at the still-running session. 
Could the stopping case get its own error message?



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, TimeoutException, 
TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepException, SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.PythonException
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.{ThreadUtils, 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 = _
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      require(active.isRunning && active.sitePackages == sitePackages,
+        "In-process Python is stopping or already initialized with different 
sitePackages")
+    } else {
+      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) {
+    private val executor = 
ThreadUtils.newDaemonSingleThreadExecutor("inprocess-python")
+    @volatile private var running = true
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+
+    def isRunning: Boolean = running
+    def isTerminated: Boolean = executor.isTerminated
+
+    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)
+        candidate.exec(
+          """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(_configured + _added))
+            |sys.path[:] = _preferred + [p for p in sys.path if p not in 
_preferred]
+            |del _site_packages, _configured, _before, _added, _preferred
+            |""".stripMargin)
+        candidate.exec("from pyspark.inprocess.runtime import " +
+          "_inprocess_invoke, _inprocess_register, _inprocess_release, _udfs")
+        interp = candidate
+      } catch {
+        case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
candidate.close() }
+      }
+    }
+
+    /** Enqueue cleanup after outstanding calls without creating an executor 
or waiting. */
+    def release(handles: Seq[String]): Unit = synchronized {
+      if (running && handles.nonEmpty) {
+        executor.submit(new Runnable {
+          override def run(): Unit = {
+            if (interp != null) interp.invoke("_inprocess_release", 
handles.asJava)
+          }
+        })
+      }
+      // During shutdown the queued close clears all remaining handles.
+    }
+
+    /** A timeout bounds plugin stop, not native execution or CDI buffer 
ownership. */
+    def shutdown(waitMillis: Long = 5000L): Unit = {
+      synchronized {
+        if (running) {
+          running = false
+          executor.submit(new Runnable {
+            override def run(): Unit = {
+              if (interp != null) {
+                try {
+                  interp.exec("_udfs.clear()")
+                } finally {
+                  try { interp.close() } finally { interp = null }
+                }
+              }
+            }
+          })
+          executor.shutdown()
+        }
+      }
+      try {
+        if (!executor.awaitTermination(waitMillis, TimeUnit.MILLISECONDS)) {
+          logWarning("In-process Python is still stopping; native work and its 
buffers " +
+            "remain alive until the invocation finishes or the process exits.")
+        }
+      } catch {
+        case _: InterruptedException => Thread.currentThread().interrupt()
+      }
+    }
+
+    def register(
+        handle: String,
+        serializedUdf: Array[Byte],
+        returnTypeJson: String,
+        timeZoneId: String,
+        pythonVersion: String,
+        largeVarTypes: Boolean): Unit = {
+      // Bulk-copy on the task thread. JEP's PyJBuffer supports memoryview 
without per-byte JNI.
+      val command = ByteBuffer.allocateDirect(serializedUdf.length)
+      command.put(serializedUdf).flip()
+      onInterpreterThread {
+        withPythonException {
+          interp.invoke("_inprocess_register", handle, command, 
returnTypeJson, timeZoneId,
+            pythonVersion, java.lang.Boolean.valueOf(largeVarTypes))
+        }
+      }
+    }
+
+    def invoke(
+        handle: String,
+        inputArrayPtrs: Array[Long],
+        inputSchemaPtrs: Array[Long],
+        outputArrayAddr: Long,
+        outputSchemaAddr: Long,
+        expectedRows: Int,
+        argumentNames: Array[String]): Long = onInterpreterThread {
+      val start = System.nanoTime()
+      val arrayPtrs = inputArrayPtrs.map(java.lang.Long.valueOf).toSeq.asJava
+      val schemaPtrs = inputSchemaPtrs.map(java.lang.Long.valueOf).toSeq.asJava
+      withPythonException {
+        interp.invoke("_inprocess_invoke", handle, arrayPtrs, schemaPtrs,
+          java.lang.Long.valueOf(outputArrayAddr), 
java.lang.Long.valueOf(outputSchemaAddr),
+          java.lang.Integer.valueOf(expectedRows), argumentNames.toSeq.asJava)
+      }
+      (System.nanoTime() - start) / 1000000

Review Comment:
   Truncating to milliseconds on each call means every batch that takes less 
than 1 ms adds 0. `pythonProcessingTime` can therefore show 0 ms for a long 
query made of fast batches. `pythonInitTime` in the factory is truncated the 
same way. Accumulating nanoseconds and converting once would avoid this.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, TimeoutException, 
TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepException, SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.PythonException
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.{ThreadUtils, 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 = _
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      require(active.isRunning && active.sitePackages == sitePackages,
+        "In-process Python is stopping or already initialized with different 
sitePackages")
+    } else {
+      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) {
+    private val executor = 
ThreadUtils.newDaemonSingleThreadExecutor("inprocess-python")
+    @volatile private var running = true
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+
+    def isRunning: Boolean = running
+    def isTerminated: Boolean = executor.isTerminated
+
+    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)
+        candidate.exec(

Review Comment:
   `register` and `invoke` convert `BaseException` to `RuntimeError` because an 
uncaught `SystemExit` makes JEP call `exit()`. These two `exec` calls have no 
such guard. A `.pth` line in `sitePackages` (`site.addpackage` only catches 
`Exception`) or an import-time `sys.exit()` would terminate the executor JVM, 
or the driver in local mode. The exit code can be 0, and the plugin's error log 
never runs. Wrapping both scripts in `try/except BaseException` would make this 
consistent.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,258 @@
+/*
+ * 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.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, TimeoutException, 
TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepException, SharedInterpreter}
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.PythonException
+import org.apache.spark.internal.Logging
+import org.apache.spark.util.{ThreadUtils, 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 = _
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      require(active.isRunning && active.sitePackages == sitePackages,
+        "In-process Python is stopping or already initialized with different 
sitePackages")
+    } else {
+      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) {
+    private val executor = 
ThreadUtils.newDaemonSingleThreadExecutor("inprocess-python")
+    @volatile private var running = true
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+
+    def isRunning: Boolean = running
+    def isTerminated: Boolean = executor.isTerminated
+
+    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 {

Review Comment:
   Once work has started, cancellation is ignored, and all 
register/invoke/release calls from every task on the executor go through this 
single thread. So one UDF that never returns blocks all in-process UDFs on that 
executor indefinitely, across jobs and sessions. The worker path can kill a 
stuck worker via `spark.python.task.killTimeout`; there is no equivalent here. 
The docs only describe the per-task effect. Could we at least document the 
executor-wide impact, or consider a watchdog that fails the executor after a 
timeout?



##########
python/pyspark/inprocess/udf.py:
##########
@@ -0,0 +1,183 @@
+#
+# 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 inspect import getfullargspec
+from typing import Any, Callable, Optional, Union
+
+from pyspark import Accumulator, Broadcast, cloudpickle
+from pyspark.errors import PySparkValueError
+from pyspark.sql.column import Column
+from pyspark.sql.types import DataType
+
+
+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) -> bytes:
+    buffer = io.BytesIO()
+    _InProcessPickler(buffer).dump(func)
+    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: DataType, deterministic: 
bool = True) -> None:
+        self._return_type: DataType = return_type
+        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
+
+    def _serialize(self) -> bytes:
+        if self._serialized is None:
+            self._serialized = _serialize_udf(self._func)
+        return self._serialized
+
+    def __call__(self, *cols: Union[Column, str], **kwargs: Union[Column, 
str]) -> Column:
+        """
+        Create a ``Column`` expression invoking this UDF with the given 
columns.
+
+        Args:
+            *cols: Spark ``Column`` objects (e.g. ``df.value``, ``col("x")``)
+
+        Returns:
+            pyspark.sql.Column
+        """
+        from pyspark import SparkContext
+        from pyspark.sql.classic.column import _to_java_column
+
+        sc = SparkContext._active_spark_context
+        if sc is None:
+            raise RuntimeError(
+                "No active SparkContext. Start a SparkSession before calling 
an inprocess_udf."
+            )
+
+        jvm = sc._jvm
+        assert jvm is not None
+
+        # Convert Python Column objects to JVM Column objects
+        if not cols and not kwargs:
+            raise PySparkValueError(
+                errorClass="INVALID_PANDAS_UDF",
+                messageParameters={"detail": "An inprocess_udf requires at 
least one argument."},
+            )
+        jcols = [_to_java_column(c) for c in cols]
+        jcols.extend(
+            jvm.PythonSQLUtils.namedArgumentExpression(name, 
_to_java_column(value))
+            for name, value in kwargs.items()
+        )
+
+        # Build a Java ArrayList (py4j vararg spread doesn't work with 
Arrays.asList)
+        jlist = jvm.java.util.ArrayList()
+        for jcol in jcols:
+            jlist.add(jcol)
+
+        # Use the existing PythonUDF planning contracts with an in-process 
eval type.
+        jcol = 
jvm.org.apache.spark.sql.execution.python.InProcessPythonUDFBuilder.build(
+            self._name,
+            self._serialize(),
+            self._return_type.json(),

Review Comment:
   A DDL string return type such as `@inprocess_udf("long")` is accepted by 
`udf`, `pandas_udf`, and `arrow_udf`. Here it passes decoration and only fails 
at this line with `AttributeError: 'str' object has no attribute 'json'`. Could 
we parse it with `_parse_datatype_string`, or reject non-`DataType` values 
early? Also, the wrapper does not expose `returnType`, `evalType`, 
`asNondeterministic`, or `__name__`/`__doc__` the way other PySpark UDF objects 
do.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,196 @@
+#
+# 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
+import traceback as _traceback
+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.sql.types import _parse_datatype_json_string
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType]] = {}
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    return_type_json: str,
+    timezone: str,
+    python_version: str,
+    large_var_types: 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.
+        func = cloudpickle.loads(memoryview(serialized_udf))
+        expected_type = to_arrow_type(
+            _parse_datatype_json_string(return_type_json),
+            timezone=timezone,
+            prefers_large_types=large_var_types,
+            error_on_duplicated_field_names_in_struct=True,
+        )
+        _udfs[handle] = (func, expected_type)
+    except BaseException:
+        # In JEP, an uncaught SystemExit can terminate the entire executor JVM.
+        raise RuntimeError(_UDF_TRACEBACK_SENTINEL + _traceback.format_exc()) 
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 _check_nested_nulls(array: pa.Array, expected_type: pa.DataType) -> None:
+    def check_field(values: pa.Array, field: pa.Field) -> None:
+        if not field.nullable and values.null_count:
+            raise ValueError(f"In-process UDF returned nulls in non-nullable 
field {field.name}")
+        _check_nested_nulls(values, field.type)
+
+    if pa.types.is_struct(expected_type):
+        # Children under a null parent do not contribute values to the result.
+        visible = pc.filter(array, pc.is_valid(array))

Review Comment:
   This `pc.filter`, and the `concat_arrays(filter(...))` for maps, copies the 
whole result on every batch at each nesting level. It does so even when the 
expected type has no non-nullable descendants, in which case `check_field` can 
never raise. Null map keys are already rejected by `result.validate()`. 
Precomputing at registration whether any non-nullable descendant exists, and 
skipping the filter when `null_count == 0`, would keep nested results closer to 
zero-copy.



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