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


##########
python/pyspark/inprocess/udf.py:
##########
@@ -0,0 +1,251 @@
+#
+# 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 signature
+from typing import Any, Callable, Optional, Union
+
+from pyspark import Accumulator, Broadcast, cloudpickle
+from pyspark.errors import PySparkNotImplementedError, 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 PySparkNotImplementedError(
+                errorClass="NOT_IMPLEMENTED",
+                messageParameters={
+                    "feature": "Spark broadcasts or accumulators in in-process 
UDFs"
+                },
+            )
+        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: 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")
+
+        from pyspark.sql.pandas.utils import require_minimum_pyarrow_version
+
+        require_minimum_pyarrow_version()
+        if not signature(func).parameters:
+            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)
+            from pyspark.sql.pandas.types import to_arrow_type
+
+            to_arrow_type(parsed, timezone="UTC", 
error_on_duplicated_field_names_in_struct=True)
+            self._parsed_return_type = parsed
+        return self._parsed_return_type
+
+    @property
+    def deterministic(self) -> bool:
+        return self._deterministic
+
+    def asNondeterministic(self) -> "InProcessUDFWrapper":
+        self._deterministic = False
+        return self
+
+    def _serialize(self) -> bytes:
+        if self._serialized is None:
+            # Validate before caching the command, including driver-only UDT 
definitions.
+            self.returnType
+            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.sql.classic.column import _to_java_column
+        from pyspark.sql.utils import get_active_spark_context, is_remote
+
+        if is_remote():
+            raise PySparkNotImplementedError(
+                errorClass="NOT_IMPLEMENTED",
+                messageParameters={"feature": "In-process Python UDFs in Spark 
Connect"},
+            )
+        sc = get_active_spark_context()
+
+        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(

Review Comment:
   **[Medium] Each call sends the whole pickled closure to the JVM again, and 
each plan keeps its own copy.**
   
   `_serialize()` caches the bytes on the Python side, but every `__call__` 
passes them to `InProcessPythonUDFBuilder.build` through Py4J, which sends 
`bytes` base64-encoded, and `build` wraps the new `byte[]` in a new 
`SimplePythonFunction`. Worker UDFs differ in two ways:
   
   - `UserDefinedFunction._judf` creates the JVM function once, and later calls 
pass only the columns (`judf.apply(...)`).
   - `_prepare_for_python_RDD` (`core/rdd.py` L5111) replaces a command larger 
than `spark.broadcast.UDFCompressionThreshold` (1 MiB by default) with a 
broadcast, so the expression holds only a small reference.
   
   Since broadcasts are rejected here, large state such as a model has to live 
in the closure, so large commands are more likely on this path than on the 
worker path. For example, with a 100 MB closure, `for c in cols: df = 
df.withColumn(c, f(df[c]))` over 50 columns sends about 6.7 GB through Py4J, 
and the resulting plan holds 50 distinct 100 MB arrays in the driver heap. 
Catalyst also hashes and compares these arrays in full when it compares the 
expressions (`semanticEquals` compares contents now), and every stage that uses 
the UDF carries the whole closure in its task binary.
   
   Suggestion: create the JVM-side function once per wrapper, as `_judf` does, 
and pass only the columns per call. Whether to broadcast large commands could 
be a separate decision, since the embedded runtime would have to read the 
broadcast without a worker.



##########
python/pyspark/inprocess/udf.py:
##########
@@ -0,0 +1,251 @@
+#
+# 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 signature
+from typing import Any, Callable, Optional, Union
+
+from pyspark import Accumulator, Broadcast, cloudpickle
+from pyspark.errors import PySparkNotImplementedError, 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 PySparkNotImplementedError(
+                errorClass="NOT_IMPLEMENTED",
+                messageParameters={
+                    "feature": "Spark broadcasts or accumulators in in-process 
UDFs"
+                },
+            )
+        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: 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")
+
+        from pyspark.sql.pandas.utils import require_minimum_pyarrow_version
+
+        require_minimum_pyarrow_version()
+        if not signature(func).parameters:
+            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)
+            from pyspark.sql.pandas.types import to_arrow_type
+
+            to_arrow_type(parsed, timezone="UTC", 
error_on_duplicated_field_names_in_struct=True)
+            self._parsed_return_type = parsed
+        return self._parsed_return_type
+
+    @property
+    def deterministic(self) -> bool:
+        return self._deterministic
+
+    def asNondeterministic(self) -> "InProcessUDFWrapper":
+        self._deterministic = False
+        return self
+
+    def _serialize(self) -> bytes:
+        if self._serialized is None:
+            # Validate before caching the command, including driver-only UDT 
definitions.
+            self.returnType
+            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.sql.classic.column import _to_java_column
+        from pyspark.sql.utils import get_active_spark_context, is_remote
+
+        if is_remote():
+            raise PySparkNotImplementedError(
+                errorClass="NOT_IMPLEMENTED",
+                messageParameters={"feature": "In-process Python UDFs in Spark 
Connect"},
+            )
+        sc = get_active_spark_context()
+
+        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.returnType.json(),
+            jlist,
+            self._deterministic,
+            "%d.%d" % sys.version_info[:2],
+        )
+
+        return Column(jcol)
+
+
+def inprocess_udf(return_type: Union[DataType, str], deterministic: bool = 
True) -> Callable:

Review Comment:
   **[Low] `return_type` differs from the `returnType` keyword of `udf`, 
`pandas_udf` and `arrow_udf`.**
   
   All three existing factories take `returnType` (`udf(f=None, 
returnType=StringType(), *, useArrow=None)`, `pandas_udf(f=None, 
returnType=None, functionType=None)` and `arrow_udf(f=None, returnType=None, 
functionType=None)`), and this wrapper's own property is `returnType` too. The 
guide (L624) says that `inprocess_udf` and `pandas_udf` "have nearly identical 
call-site syntax", but rewriting `@pandas_udf(returnType=LongType())` as 
`@inprocess_udf(returnType=LongType())` raises `TypeError: inprocess_udf() got 
an unexpected keyword argument 'returnType'`.
   
   Since this is a new public API (`versionadded:: 4.4.0`), the keyword is hard 
to change after the release. Could it take `returnType` instead, with the 
docstring and guide examples updated accordingly?



##########
sql/core/src/test/scala/org/apache/spark/sql/execution/python/InProcessEvaluatorTestUtils.scala:
##########
@@ -0,0 +1,114 @@
+/*
+ * 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.concurrent.{CountDownLatch, TimeUnit}
+import java.util.concurrent.atomic.AtomicInteger
+
+import org.apache.spark.TaskContext
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{AttributeReference, 
GenericInternalRow, UnsafeProjection}
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.types.{DataType, LongType, StructField, StructType}
+
+/** Fixtures for in-process evaluator tests that need no Python. */
+private[python] object InProcessEvaluatorTestUtils {
+
+  /** Every metric that an evaluator may update, as `PythonSQLMetrics` defines 
them. */
+  def allMetrics(): Map[String, SQLMetric] =
+    (PythonSQLMetrics.pythonSizeMetricsDesc ++ 
PythonSQLMetrics.pythonTimingMetricsDesc ++
+      PythonSQLMetrics.pythonOtherMetricsDesc).keys.map(_ -> new 
SQLMetric("sum", 0L)).toMap
+
+  def thread(body: => Unit): Thread = {
+    val t = new Thread(() => body)
+    t.start()
+    t
+  }
+
+  /**
+   * An evaluator without UDFs over `rowCount` rows of one long column, so 
that its iterator
+   * runs without Python. The input blocks on `gate` before it reads row 
`blockAt`: in
+   * `hasNext`, or in `next` if `blockInNext`. If `blockInCopy`, it returns 
that row instead,
+   * which blocks when its value is read, i.e. when the evaluator copies it 
into a batch. For
+   * `ReadBack`, which takes any row, `copied` counts the copies of input rows.
+   */
+  class BlockingInput(
+      joinInput: InProcessArrowEvalPythonEvaluatorFactory.JoinInput,
+      val context: TaskContext,
+      session: => InProcessPythonRuntime.InterpreterSession,
+      rowCount: Int = Int.MaxValue,
+      blockAt: Int = -1,
+      blockInNext: Boolean = false,
+      blockInCopy: Boolean = false,
+      batchSize: Int = 10) {
+    val reached = new CountDownLatch(1)
+    val gate = new CountDownLatch(1)
+    val pulled = new AtomicInteger()
+    val copied = new AtomicInteger()
+    private val column = AttributeReference("x", LongType)()
+    private val toUnsafe = UnsafeProjection.create(Array[DataType](LongType))
+
+    private def await(): Unit = {
+      reached.countDown()
+      gate.await(10, TimeUnit.SECONDS)
+    }
+
+    private def block(inNext: Boolean): Unit = {
+      if (!blockInCopy && inNext == blockInNext && pulled.get == blockAt) 
await()
+    }
+
+    private val rows: Iterator[InternalRow] = new Iterator[InternalRow] {
+      override def hasNext: Boolean = { block(inNext = false); pulled.get < 
rowCount }
+      override def next(): InternalRow = {
+        block(inNext = true)
+        val blocks = blockInCopy && pulled.get == blockAt
+        val value = pulled.incrementAndGet().toLong
+        if (blocks || joinInput == 
InProcessArrowEvalPythonEvaluatorFactory.ReadBack) {

Review Comment:
   **[Low, test] With `Buffered`, `blockInCopy` fails with a 
`ClassCastException` instead of blocking.**
   
   For the blocking row, this returns a `GenericInternalRow` in every mode. In 
the `Buffered` mode, `pullRow` first runs 
`queue.add(row.asInstanceOf[UnsafeRow])` 
(`InProcessArrowEvalPythonEvaluatorFactory.scala` L306), so that row throws 
`ClassCastException` before `getLong` is ever called. The scaladoc (L46-47) 
describes `blockInCopy` without restricting it to `ReadBack`, so a `Buffered` 
counterpart of "task completion waits for a consumer copying an input row read 
back from Arrow" would fail for an unrelated reason.
   
   Since `UnsafeRow` is final and cannot block when it is read, could the 
fixture `require(!blockInCopy || joinInput == ReadBack)` and say so in the 
scaladoc?



##########
core/src/main/scala/org/apache/spark/internal/config/Python.scala:
##########
@@ -56,6 +57,29 @@ private[spark] object Python {
     .bytesConf(ByteUnit.MiB)
     .createOptional
 
+  // Defined before the config entry, whose validator captures it.
+  private[spark] val IN_PROCESS_PATH_RULE = "In-process Python site-packages 
paths cannot " +
+    "contain single quotes, newlines, NUL, surrogate characters (including 
supplementary " +
+    "Unicode characters) or the platform path separator"
+
+  val IN_PROCESS_SITE_PACKAGES = 
ConfigBuilder("spark.inprocess.python.sitePackages")

Review Comment:
   **[Low] This adds a new top-level `spark.inprocess` namespace, spelled 
differently from the SQL conf of the same feature.**
   
   The SQL conf is 
`spark.sql.execution.pythonUDF.inProcess.fullValidation.enabled` (camel case 
`inProcess`), while this one is `spark.inprocess.python.sitePackages`, with 
lower case `inprocess` as a new top-level prefix. The other settings in this 
file use `spark.python.*` (e.g. `spark.python.worker.reuse`, 
`spark.python.task.killTimeout`, `spark.python.udf.pipelined.enabled`) or 
`spark.executor.pyspark.*`.
   
   Since configuration names are public and hard to change after 4.4.0, could 
this be e.g. `spark.python.inProcess.sitePackages`?



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalPythonEvaluatorFactory.scala:
##########
@@ -0,0 +1,564 @@
+/*
+ * 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.file.Files
+import java.util.UUID
+import java.util.concurrent.TimeUnit
+import java.util.concurrent.atomic.AtomicInteger
+import java.util.concurrent.locks.ReentrantLock
+
+import scala.collection.mutable.ArrayBuffer
+import scala.jdk.CollectionConverters._
+
+import com.google.common.util.concurrent.Uninterruptibles
+import org.apache.arrow.c.{ArrowArray, ArrowSchema, BaseStruct}
+import org.apache.arrow.util.AutoCloseables
+import org.apache.arrow.vector.VectorSchemaRoot
+
+import org.apache.spark.{SparkEnv, SparkException, TaskContext}
+import org.apache.spark.api.python.ChainedPythonFunctions
+import org.apache.spark.memory.MemoryConsumer
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression, 
JoinedRow, PythonUDF, UnsafeProjection, UnsafeRow}
+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._
+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, unless all of them are UDF arguments that read back from Arrow 
unchanged.
+ * Each batch owns its Arrow buffers so Python can safely retain input arrays.
+ *
+ * The evaluator owns its queue, so that cleanup at task completion is 
coordinated with a
+ * consumer on another thread, such as a pipelined Python writer or a 
TRANSFORM feed thread.
+ */
+class InProcessArrowEvalPythonEvaluatorFactory(
+    childOutput: Seq[Attribute],
+    udfs: Seq[PythonUDF],
+    output: Seq[Attribute],
+    batchSize: Int,
+    maxBytes: Long,
+    timeZoneId: String,
+    largeVarTypes: Boolean,
+    hideTraceback: Boolean,
+    simplifiedTraceback: Boolean,
+    tracebackWithLocals: Boolean,
+    fullValidation: Boolean,
+    metrics: Map[String, SQLMetric])
+  extends EvalPythonEvaluatorFactory(childOutput, udfs, output) {
+
+  private[python] def runtimeSession: 
InProcessPythonRuntime.InterpreterSession =
+    InProcessPythonRuntime.currentSession
+
+  /** Unused: `evaluateJoined` always evaluates the UDFs. */
+  override protected def evaluate(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputSchema: StructType,
+      context: TaskContext): Iterator[InternalRow] =
+    throw SparkException.internalError("In-process UDFs are evaluated with 
their input rows")
+
+  override protected def evaluateJoined(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputs: Seq[Expression],
+      inputSchema: StructType,
+      context: TaskContext): Option[Iterator[InternalRow]] = {
+    import InProcessArrowEvalPythonEvaluatorFactory.{Buffered, ReadBack, 
readsBack}
+    val inputColumns = inputs.length == childOutput.length && 
inputs.zip(childOutput).forall {
+      case (a: Attribute, c) => a.exprId == c.exprId
+      case _ => false
+    }
+    // If all input columns are UDF arguments, they are written to Arrow 
regardless. Read them
+    // back from the exported input vectors instead of buffering every input 
row, if their
+    // values read back from Arrow exactly as written and as fast as an unsafe 
row copy.
+    val joinInput = if (inputColumns && inputSchema.forall(f => 
readsBack(f.dataType))) {
+      ReadBack
+    } else if (inputColumns) {
+      Buffered(None)
+    } else {
+      // Each projected row is written to Arrow before the next input row is 
pulled, so the
+      // arguments go into a reused buffer rather than being copied value by 
value.
+      val projection = UnsafeProjection.create(inputs, childOutput)
+      projection.initialize(context.partitionId())
+      Buffered(Some(projection))
+    }
+    Some(evaluateBatches(funcs, argMetas, rows, inputSchema, context, 
joinInput))
+  }
+
+  private[python] def evaluateBatches(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputSchema: StructType,
+      context: TaskContext,
+      joinInput: InProcessArrowEvalPythonEvaluatorFactory.JoinInput): 
Iterator[InternalRow] = {
+    import InProcessArrowEvalPythonEvaluatorFactory.{Buffered, ReadBack}
+    ArrowUtils.failDuplicatedFieldNames(inputSchema)
+    val functions = funcs.map { case (chain, _) =>
+      if (chain.funcs.size != 1) {
+        throw SparkException.internalError(
+          "In-process UDF chains must use separate evaluation nodes")
+      }
+      chain.funcs.head
+    }
+    val inputOrdinals = argMetas.map(_.map(_.offset))
+    def checkCancellation(): Unit = context.killTaskIfInterrupted()
+
+    val expectedFields = udfs.map { udf =>
+      ArrowUtils.toArrowField("result", udf.dataType, true, timeZoneId, 
largeVarTypes)
+    }
+    val processingTime = new 
InProcessArrowEvalPythonEvaluatorFactory.NanosecondTimer(
+      metrics("pythonProcessingTime"))
+    val initTime = new 
InProcessArrowEvalPythonEvaluatorFactory.NanosecondTimer(
+      metrics("pythonInitTime"))
+    val arrowSchema = ArrowUtils.toArrowSchema(inputSchema, timeZoneId, 
largeVarTypes)
+    // Capture before consuming input: an old task must never join a later 
context's session.
+    val runtime = runtimeSession
+    // Rows are copied out of the queue and Arrow vectors before they are 
returned, so they
+    // remain valid after task completion releases those, on whichever thread 
consumes them.
+    val resultProj = UnsafeProjection.create(output, output)
+    // Spill files go into a directory of the queue's own, created with the 
first disk queue, so
+    // that task completion can delete them when it cannot close the queue.
+    @volatile var spillDir: File = null
+    // Guarded by the queue's monitor.
+    var queueAbandoned = false
+    val (queue, projection) = joinInput match {
+      case Buffered(projection) =>
+        val localDir = new File(Utils.getLocalDir(SparkEnv.get.conf))
+        val serializerManager = SparkEnv.get.serializerManager
+        // Only the consumer holding the iterator's lock adds and removes rows.
+        val queue = new HybridRowQueue(context.taskMemoryManager(), localDir,
+            childOutput.length, serializerManager, lockFree = true) {
+          override protected def createDiskQueue(): RowQueue = synchronized {
+            if (spillDir == null) {
+              spillDir = Files.createTempDirectory(localDir.toPath, 
"inprocess-udf-").toFile
+            }
+            DiskRowQueue(Files.createTempFile(spillDir.toPath, "buffer", 
"").toFile,
+              childOutput.length, serializerManager)
+          }
+
+          // Once task completion leaves the queue to the executor, it must 
not spill for other
+          // consumers into a directory that nothing deletes.
+          override def spill(size: Long, trigger: MemoryConsumer): Long = 
synchronized {
+            if (queueAbandoned) 0L else super.spill(size, trigger)
+          }
+
+          // Queues of a task are distinct memory consumers, whatever their 
case-class fields.
+          override def equals(other: Any): Boolean = this eq 
other.asInstanceOf[AnyRef]
+          override def hashCode(): Int = System.identityHashCode(this)
+          override def canEqual(other: Any): Boolean = false
+        }
+        (queue, projection.orNull)
+      case ReadBack => (null, null)
+    }
+    val joined = new JoinedRow
+    val handles = functions.map(_ => UUID.randomUUID().toString)
+    var registered = false
+    var writer: ArrowWriter = null
+    val results = ArrayBuffer.empty[ArrowColumnVector]
+    var startedAt = 0L
+
+    def closeBatch(): Unit = {
+      val resources = ArrayBuffer.empty[AutoCloseable]
+      resources ++= results
+      results.clear()
+      if (writer != null) {
+        resources += writer.root
+        writer = null
+      }
+      AutoCloseables.close(resources.asJava)
+    }
+
+    val resources = new 
InProcessArrowEvalPythonEvaluatorFactory.IteratorResources(
+      // Closing the queue deletes the spill files it tracks; deleteQuietly 
also removes any
+      // other, without starting a process or throwing, also on an interrupted 
thread.
+      releaseTaskMemory = () => if (queue != null) {
+        Utils.tryWithSafeFinally(queue.close())(Utils.deleteQuietly(spillDir))
+      },
+      abandonTaskMemory = () => if (queue != null) {
+        queue.synchronized {
+          queueAbandoned = true
+          Utils.deleteQuietly(spillDir)
+        }
+      },
+      releaseOthers = () => {
+        if (startedAt != 0L) {
+          metrics("pythonTotalTime") += (System.nanoTime() - startedAt) / 
1000000
+        }
+        Utils.tryWithSafeFinally {
+          closeBatch()
+        } {
+          if (registered) runtime.release(handles)
+        }
+      })
+
+    context.addTaskCompletionListener[Unit](_ => resources.close())
+
+    new Iterator[InternalRow] {
+      private var batchIter: Iterator[InternalRow] = Iterator.empty
+
+      private def endOfInput: Nothing =
+        throw new NoSuchElementException("End of in-process UDF input")
+
+      // Releases the resources on failure without replacing its exception.
+      private def fail(t: Throwable): Nothing =
+        Utils.tryWithSafeFinally { throw t } { resources.close() }
+
+      // Called with the lock held.
+      private def hasNextLocked: Boolean = {
+        if (startedAt == 0L) startedAt = System.nanoTime()
+        checkCancellation()
+        val available = batchIter.hasNext || {
+          var open = false
+          val more = try {
+            resources.startReadingInput() && rows.hasNext
+          } finally {
+            open = resources.endReadingInput()
+          }
+          more && open
+        }
+        if (!available) resources.close()
+        available
+      }
+
+      // Each call takes the lock without allocating a closure per row. Within 
a batch, only
+      // the consumer advances `batchIter`, whose `hasNext` compares row 
indexes.
+      override def hasNext: Boolean = {
+        if (batchIter.hasNext && !resources.isClosed) return true
+        if (!resources.enter()) return false
+        val available = try {
+          try hasNextLocked catch { case t: Throwable => fail(t) }
+        } finally {
+          resources.exit()
+        }
+        // Task completion may have closed the iterator while the input was 
being read.
+        available && !resources.isClosed
+      }
+
+      override def next(): InternalRow = {
+        if (!resources.enter()) endOfInput
+        try {
+          try {
+            if (!hasNextLocked) endOfInput
+            if (!batchIter.hasNext) nextBatch()
+            val result = batchIter.next()
+            resultProj(if (queue != null) joined(queue.remove(), result) else 
result)
+          } catch {
+            case t: Throwable => fail(t)
+          }
+        } finally {
+          resources.exit()
+        }
+      }
+
+      // Runs Python without the lock unless task completion already happened, 
and ends the
+      // input instead of returning the result if it happens meanwhile.
+      private def python[T](body: => T): T = {
+        if (resources.isClosed) endOfInput
+        val result = resources.withoutLock(body)
+        if (resources.isClosed) endOfInput
+        result
+      }
+
+      /**
+       * Writes the next input row to the batch, returning false at the end of 
input or once
+       * task completion happened. If it happens while the row is read, the 
row is dropped.
+       */
+      private def pullRow(): Boolean = {
+        var open = false
+        val row = try {
+          // Checked again after `hasNext`, which may wait for input, not to 
read another row.
+          if (resources.startReadingInput() && rows.hasNext && 
!resources.isClosed) {
+            rows.next()
+          } else {
+            null
+          }
+        } finally {
+          open = resources.endReadingInput()
+        }
+        if (!open) endOfInput
+        if (row == null) return false
+        if (queue != null) queue.add(row.asInstanceOf[UnsafeRow])
+        writer.write(if (projection != null) projection(row) else row)
+        true
+      }
+
+      // Called with the lock held.
+      private def nextBatch(): Unit = {
+        closeBatch()
+        val root = VectorSchemaRoot.create(arrowSchema, 
ArrowUtils.rootAllocator)
+        writer = try {
+          ArrowWriter.create(root)
+        } catch {
+          case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
root.close() }
+        }
+        // Task completion stops the fill within a row, and Python never sees 
a partial batch.
+        var count = 0
+        while ((batchSize <= 0 || count < batchSize) &&
+            (count == 0 || writer.sizeInBytes() < maxBytes) && {
+              checkCancellation()
+              pullRow()
+            }) {
+          count += 1
+        }
+        if (resources.isClosed) endOfInput
+        if (!registered) {
+          // Mark before registering so failure after any registration still 
cleans up.
+          registered = true
+          functions.indices.foreach { i =>
+            val func = functions(i)
+            initTime.add(python(runtime.register(handles(i), 
func.command.toArray,
+              expectedFields(i), func.pythonVer, hideTraceback, 
simplifiedTraceback,
+              tracebackWithLocals, fullValidation)))
+          }
+        }
+        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 track[S <: BaseStruct](struct: S): S = {
+            val closer: AutoCloseable = () => 
InProcessArrowBridge.closeStruct(struct)
+            structs += closer
+            struct
+          }
+          def array(): ArrowArray = 
track(ArrowArray.allocateNew(ArrowUtils.rootAllocator))
+          def schema(): ArrowSchema = 
track(ArrowSchema.allocateNew(ArrowUtils.rootAllocator))
+          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))
+            }
+            processingTime.add(python(runtime.invoke(
+              handle,
+              inArrays.map(_.memoryAddress()).toArray,
+              inSchemas.map(_.memoryAddress()).toArray,
+              outArray.memoryAddress(), outSchema.memoryAddress(),
+              count, argMetas(udfIndex).map(_.name.getOrElse("")))))
+            results += InProcessArrowBridge.cdiToColumn(
+              outArray, outSchema, Some(expectedFields(udfIndex)))
+            metrics("pythonDataReceived") += 
results.last.getValueVector.getBufferSize
+          } {
+            AutoCloseables.close(structs.asJava)
+          }
+        }
+
+        metrics("pythonNumRowsReceived") += count
+        // Input vectors are closed with the writer's root, not with the 
results.
+        val inputs = if (joinInput == ReadBack) {
+          writer.root.getFieldVectors.asScala.map(new ArrowColumnVector(_))
+        } else {
+          Nil
+        }
+        val columns = (inputs ++ results).toArray[ColumnVector]
+        batchIter = new ColumnarBatch(columns, count).rowIterator().asScala
+      }
+    }
+  }
+}
+
+private[python] object InProcessArrowEvalPythonEvaluatorFactory {
+  /** How the evaluator joins input rows with their results. */
+  sealed trait JoinInput
+  /** Read the input columns back from the exported Arrow input vectors. */
+  case object ReadBack extends JoinInput
+  /** Buffer the input rows, writing their arguments, projected if needed, to 
Arrow. */
+  case class Buffered(projection: Option[UnsafeProjection]) extends JoinInput
+
+  /**
+   * Whether `ArrowColumnVector` returns exactly the values `ArrowWriter` 
wrote for this type,
+   * and an unsafe projection copies them about as fast as an unsafe row. 
Types with derived
+   * Arrow representations, such as intervals, nanosecond timestamps, TIME, 
Variant, geospatial
+   * types and UDTs, keep the original rows instead. So do arrays and maps, 
which a projection
+   * copies element by element out of Arrow, but with a single copy out of an 
unsafe row, and
+   * decimals, which Arrow reads back through a `BigDecimal` per value.
+   */
+  def readsBack(dataType: DataType): Boolean = dataType match {
+    case NullType | BooleanType | ByteType | ShortType | IntegerType | 
LongType |
+        FloatType | DoubleType | BinaryType | DateType | TimestampType | 
TimestampNTZType => true
+    case _: StringType => true
+    case StructType(fields) => fields.forall(f => readsBack(f.dataType))
+    case _ => false
+  }
+
+  /**
+   * Coordinates cleanup at task completion with the consumer of the 
evaluator's iterator. The
+   * consumer can run on another thread, e.g. a pipelined Python writer or a 
TRANSFORM feed
+   * thread, and the completion listener cannot tell, since a lazily computing 
parent (such as
+   * `coalesce`) can create the iterator on that thread too.
+   *
+   * The consumer holds the lock while it reads input, the row queue or Arrow 
vectors, and
+   * releases it only while this evaluator's Python runs. The listener 
(`close`) first requests
+   * closing, which the consumer checks after each input row, so the listener 
waits for at most
+   * one row before it releases task memory (the row queue), ahead of the 
executor. It releases
+   * the other resources (Arrow vectors and Python handles) too, unless Python 
is running; then
+   * the consumer releases them when Python returns.
+   *
+   * Reading one row can take long: the input can be another in-process 
evaluator, whose next
+   * row may need a batch of Python, or an upstream operator that only a later 
listener
+   * unblocks. So while the consumer reads input, the listener waits for the 
lock only
+   * briefly. Then it leaves the task memory to the executor, deleting what 
lives outside it,
+   * and the consumer releases the other resources once its row returns, 
without touching the
+   * task memory again. Otherwise the consumer may use the task memory, e.g. 
the queue, and
+   * the listener waits for the lock until it is done. This holds also when 
the input is read
+   * back from Arrow, without a queue: the consumer then copies a row that may 
point into the
+   * task memory of an upstream operator, e.g. a page of a sorter, which the 
executor frees.
+   */
+  class IteratorResources(

Review Comment:
   **[Design] The consumer threads that outlive task completion are the root 
cause, and fixing them there would cover every upstream operator.**
   
   `IteratorResources` (the lock, the read-mark protocol and the abandon path), 
the spill override after abandonment, the identity `equals` and the queue's own 
spill directory all exist because this iterator can be consumed on a thread 
that keeps running after task completion. Those threads belong to two consumers:
   
   - The pipelined runner's listener (`PythonRunner.scala` L553-561) calls 
`writerFuture.cancel(true)` and then `get()`, to wait for the writer as its 
SPARK-33277 comment says. However, `FutureTask.get()` throws 
`CancellationException` as soon as the future is cancelled, without waiting for 
`run()` to return, so the listener does not wait for the writer.
   - The feed thread of `BaseScriptTransformationExec` is a daemon thread that 
no listener joins.
   
   So any other upstream operator under these consumers, e.g. a sorter or a 
hash aggregate, can still have its task memory read after the executor frees 
it, while the protocol here protects only this evaluator's queue and row copy. 
A bounded stop-and-join in the listeners of those two consumers (bounded for 
the case where the input waits for a later listener, as handled here) would fix 
it for all operators, and could simplify this class to the fallback path.
   
   This does not have to be fixed in this PR, but it might deserve a separate 
JIRA, since the protocol here is fairly complex for one operator.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,419 @@
+/*
+ * 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.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+import java.util.concurrent.atomic.AtomicInteger
+
+import scala.collection.mutable
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, 
NamingConventionClassEnquirer, PyConfig, SharedInterpreter}
+import org.apache.arrow.c.{ArrowSchema, Data}
+import org.apache.arrow.vector.types.pojo.Field
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.{PythonException, PythonUtils}
+import org.apache.spark.internal.Logging
+import org.apache.spark.internal.config.Python
+import org.apache.spark.sql.util.ArrowUtils
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var mainConfigured = false
+  @volatile private var sharedConfigured = false
+  @volatile private var bootstrappedSitePackages: Option[Seq[String]] = None
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  // Keep JEP references out of the singleton's verifier so currentSession can 
report
+  // an uninitialized runtime even when the provided JEP JAR is absent.
+  private[python] object InterpreterConfiguration {
+    def configure(sitePackages: Seq[String]): Unit = {
+      if (!mainConfigured) {
+        // 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(
+          
PyConfig.isolated().setUseEnvironment(false).setHashSeed(0).setUseHashSeed(true))

Review Comment:
   **[Low] The isolated config also drops the user site-packages directory, 
which Python workers search.**
   
   `PyConfig_InitIsolatedConfig` sets `user_site_directory = 0`, along with 
`isolated` and `use_environment`, so `site` sets `ENABLE_USER_SITE = False` and 
never adds `site.getusersitepackages()` (e.g. 
`~/.local/lib/python3.12/site-packages`) to `sys.path`. Python workers start 
without `-s` or `-I`, so they do import packages installed with `pip install 
--user`.
   
   So on an executor (or a local driver) whose PyArrow comes from `pip install 
--user`, `arrow_udf` works, but the plugin fails at startup in the bootstrap's 
`require_minimum_pyarrow_version()` with `PACKAGE_NOT_INSTALLED` ("PyArrow >= 
18.0.0 must be installed"), wrapped in the installation checklist. The guide 
(L561-562) also says that `sitePackages` is not needed on "Executors where all 
required packages are pre-installed on the system Python path", which reads as 
if this case were covered.
   
   Suggestion: either document that the user site-packages directory is not 
searched, and that it can be listed in `spark.inprocess.python.sitePackages`, 
or add `site.getusersitepackages()` after the configured paths in the bootstrap 
when it exists, to match workers.



##########
docs/sql-pyspark-inprocess-udf.md:
##########
@@ -0,0 +1,728 @@
+---
+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 or a task still owns exported results, cleanup waits for that 
task to release
+its CDI references; the memory remains live until cleanup completes 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-aware timestamps 
are relabeled
+to `spark.sql.session.timeZone` without changing their UTC instants or copying 
their buffers.
+Timezone-naive and timezone-aware timestamps are not interchangeable. String 
and binary
+offset widths, including nested values, are converted as needed to match
+`spark.sql.execution.arrow.useLargeVarTypes`. Large, fixed-size and 
dictionary-encoded
+representations of the declared types (`large_list`, `fixed_size_list`, 
`string_view`,
+`binary_view`, `fixed_size_binary` and dictionary arrays) are cast to the 
declared type.
+These conversions can allocate new buffers. Other value types must match 
exactly: use an
+explicit PyArrow cast for numeric conversions.
+Map `keys_sorted` metadata is normalized to Spark's declared map type.
+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. Zero-length levels without a usable offsets buffer, which Arrow permits, 
are given
+one. Compatible results retain zero-copy transfer.
+Before exporting a result, the runtime performs full Arrow validation, 
including interior
+offsets, because the JVM reads result buffers without bounds checks: a 
malformed result,
+such as one built from raw buffers, could otherwise produce wrong values or 
crash the
+executor. It does not validate UTF-8 in string results, because Spark strings 
may contain
+invalid UTF-8 (for example, `CAST(X'FF' AS STRING)`). Worker-based Arrow UDFs 
do not
+validate their results. To skip the full validation, set
+`spark.sql.execution.pythonUDF.inProcess.fullValidation.enabled` to `false`; 
Arrow's
+constant-time validation and the conversions above still apply.
+
+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 argument 
extraction
+and partition-evaluator path, while its evaluator buffers and joins input rows 
itself.
+Ordinary Python UDFs continue to use Python workers.
+
+`maxRecordsPerBatch <= 0` means no row-count limit. The independent
+`spark.sql.execution.arrow.maxBytesPerBatch` limit always applies.
+Only UDF arguments are converted to Arrow. Other columns stay in Spark rows,
+buffered in a spillable queue until the results are joined back. When every 
input
+column is a UDF argument and its type, other than a decimal, an array or a 
map, reads back
+from Arrow unchanged, the output reads those columns from the Arrow input 
vectors instead
+of buffering the rows.
+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. The 
runtime retains
+each exported result until the next invocation for that task or task cleanup, 
after the JVM
+has released its references. The runtime drops its Python references on the 
interpreter
+thread, so releasing JVM results does not trigger Python finalizers on Spark 
task threads.
+Cleanup can remain queued behind another task's invocation. The rows can also 
be consumed
+on another thread, such as a pipelined Python worker's writer. Task completion 
then stops
+that consumer after the input row it is reading, and waits for it, but not for 
this
+operator's Python: it releases the buffered rows at once, and the Arrow 
vectors when Python
+returns. Reading one row can take longer when the input is another in-process 
UDF, whose
+next row may need a batch of Python, or a blocked upstream operator. While the 
consumer
+reads its input, task completion waits for at most one second, and then leaves 
the buffered
+rows to the executor and deletes their spill files. For a task that otherwise 
succeeds, the
+executor then logs "Managed memory leak detected", or fails the task if
+`spark.unsafe.exceptionOnMemoryLeak` is `true`. Each row is copied before it 
is returned, so
+it remains valid after the task releases them.
+
+UDF deserialization uses PySpark's bundled cloudpickle. Each task registers its
+own function instance once and passes a small handle for subsequent batches.
+Exception text escapes NUL, surrogates and non-BMP characters for JEP's JNI 
exception
+transport. Other characters, including non-English BMP text, remain readable.
+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 and the process 
`PYTHONPATH`
+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.
+
+The embedded interpreter uses isolated Python initialization and ignores 
Python startup
+flags from the environment (including `PYTHONFAULTHANDLER` and 
`PYTHONDEVMODE`). Spark
+explicitly restores its Python distribution paths and the executor process 
`PYTHONPATH`
+before configured `sitePackages` paths. This also supports YARN's localized 
Python archives
+when `SPARK_HOME` is absent. Python modules must not install process-wide 
signal handlers
+that replace the JVM's handlers. JEP's automatic Java package discovery is 
disabled so
+Java packages do not shadow Python packages.
+
+Executors must start with a UTF-8 locale, for example `LC_ALL=C.UTF-8` on 
systems that
+provide it. Isolated initialization ignores `PYTHONUTF8` and 
`PYTHONIOENCODING` and does
+not coerce an ASCII locale; the runtime warns if it detects one. Standard 
output and error
+use line buffering and are flushed during orderly interpreter shutdown.
+
+Spark broadcasts, accumulators, `SparkFiles`, `--py-files`, 
`spark.submit.pyFiles`,
+`SparkContext.addPyFile`, `SparkSession.addArtifacts(..., pyfile=True)`, and 
Python
+`TaskContext` are not supported by this embedded runtime. Python files may 
incidentally
+be importable on YARN through its process `PYTHONPATH`; this is not portable 
support for
+these APIs. For archived environments, access files by their configured 
executor paths,
+rather than `SparkFiles.get`. 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).
+At query execution on the driver, in-process UDFs reject positive
+`spark.executor.pyspark.memory`, `spark.sql.pyspark.udf.profiler`,
+`spark.python.profile`, `spark.python.profile.memory`, and 
`spark.pythonWorkerEnv.*`
+settings. A Python memory value of `0` means no separate limit
+and is accepted. Python runs inside the JVM, so a separate Python process
+memory limit cannot be applied. Use executor memory settings for sizing, and 
worker-based
+UDFs when these Python worker features are needed. Worker-specific logging, 
faulthandler,
+traceback-dump timers, process reuse, idle timeouts, and pipelined worker 
transport settings
+do not apply to this mode. In particular, 
`spark.sql.pyspark.worker.logging.enabled` and
+`spark.sql.execution.pyspark.udf.faulthandler.enabled` do not enable these 
worker facilities
+inside the JVM. Ordinary UDFs in the same application retain their worker 
settings.
+
+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, and DataFrame API calls report that Connect is 
unsupported. 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. Set `spark.pyspark.python` to the venv executable if ordinary worker 
UDFs in the same
+application should also use that environment. This does not select JEP's 
embedded CPython;
+JEP must be built against the intended Python version. In-process UDFs do not 
fall back to workers.
+
+```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.pyspark.python=./myvenv/bin/python3 \
+  --conf 
spark.executor.extraLibraryPath=./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.extraJavaOptions=-Djava.library.path=/opt/venv/lib/python3.11/site-packages/jep
 \
+  --conf 
spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \
+  --conf spark.executor.cores=1 \
+  --conf spark.task.cpus=1 \
+  my_app.py
+```
+
+`spark.executor.extraJavaOptions=-Djava.library.path=...` locates JEP even on 
Kubernetes
+images whose entrypoint does not propagate `spark.executor.extraLibraryPath`.
+It replaces the entire JVM native library search path; append any other 
required native
+library directories, separated by `:` on Linux.
+
+#### Option B: `--archives` with remote file upload
+
+If you cannot build a custom image, Spark on Kubernetes can distribute 
archives via a remote
+staging area (e.g. S3 or GCS). Set `spark.kubernetes.file.upload.path` to an 
object storage
+path that both the driver and executors can access. Provision JEP and Arrow 
CDI JARs
+in `/opt/inprocess/jars` on every executor using an image layer or a mounted 
volume
+before the JVM starts. The image must also have Python 3.11+ and the shared
+library matching the archived JEP build. Use the same JEP version as the 
archived
+venv. A classpath wildcard pointing inside the archive is insufficient here: 
the JVM expands wildcards
+before Spark downloads and extracts the archive.
+
+```bash
+spark-submit \
+  --master k8s://https://<k8s-api-server>:<port> \
+  --deploy-mode cluster \
+  --conf 
spark.kubernetes.container.image=my-registry/spark-python311:tested-build \
+  --conf spark.kubernetes.file.upload.path=s3a://my-bucket/spark-uploads \
+  --archives myvenv.zip#myvenv \
+  --conf 'spark.executor.extraClassPath=/opt/inprocess/jars/*' \
+  --conf 
spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \
+  --conf spark.executor.cores=1 \
+  --conf spark.task.cpus=1 \
+  --conf spark.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
+```
+
+---
+
+## Configuration Reference
+
+### `spark.plugins`
+
+| Default | `(none)` |
+|---|---|
+| **Required value** | 
`org.apache.spark.sql.execution.python.InProcessPythonPlugin` |
+
+Registers the in-process Python plugin. This initializes the 
`SharedInterpreter` on each
+executor at startup. Without this plugin, the driver rejects in-process UDF 
execution
+during physical planning, including `explain()`, before submitting tasks. 
Missing native
+dependencies are reported during plugin initialization.
+Task calls and cleanup never create or restart an interpreter.
+
+---
+
+### `spark.inprocess.python.sitePackages`
+
+| Default | `(none)` |
+|---|---|
+| **Type** | Comma-separated list of absolute or relative directory paths |
+
+Site-package directories supplied to JEP before interpreter construction, so 
the
+`jep` package must be directly importable from these directories. After 
construction,
+paths are made absolute and processed with `site.addsitedir`, including `.pth` 
files.
+Spark distribution paths and process `PYTHONPATH` take precedence over these 
directories;
+configured paths then take precedence over system paths for modules not yet 
imported.
+`.pth` files are processed too late to locate JEP itself during construction.
+
+**When you need this:** When you distribute a Python virtual environment via 
`--archives` and
+need packages from that venv to be importable inside UDFs. The problem is that 
the jep
+interpreter starts with the *system* Python's `sys.path`, which does not 
include the distributed
+venv's site-packages. Setting this config tells the plugin where to find the 
venv's packages.
+
+**Typical usage with `--archives`:**
+
+```
+spark.inprocess.python.sitePackages = ./myvenv/lib/python3.11/site-packages
+```
+
+The relative path `./myvenv/` resolves to the directory where Spark extracted 
your archive on
+the executor node. Initial archives are localized or unpacked before executor 
plugin
+initialization, so their site-packages directories are available when JEP 
starts.
+
+Paths cannot contain a single quote, newline, NUL, surrogate character 
(including
+supplementary Unicode characters), comma, or the platform path separator (`:` 
on Linux/macOS).
+Backslashes are supported. Commas separate configuration entries.
+
+The configured directories are fixed for the JVM lifetime once JEP has been 
configured,
+including when a subsequent Python bootstrap fails. Stopping a SparkContext 
does not clear
+Python's cached modules or previous `.pth` entries. Restart the executor 
process (or the
+local driver process) before changing this configuration or replacing its 
installed packages.
+
+**Multiple paths** (comma-separated):
+
+```
+spark.inprocess.python.sitePackages = 
./venv/lib/python3.11/site-packages,/opt/custom/lib
+```
+
+**When you do NOT need this:**
+- Executors where all required packages are pre-installed on the system Python 
path.
+
+---
+
+### Native library paths
+
+jep requires its native library (`libjep.so` on Linux, `libjep.dylib` on 
macOS) to be on the
+JVM's native library path. **This must be set before the JVM starts** — 
`System.setProperty()`
+has no effect after JVM startup, so runtime configuration is not possible.
+
+For YARN and Standalone, use `spark.executor.extraLibraryPath` to prepend 
JEP's directory
+to the native library search path while preserving existing directories:
+
+```
+spark.executor.extraLibraryPath = ./myvenv/lib/python3.11/site-packages/jep
+```
+
+For Kubernetes images whose entrypoint does not propagate this setting, set
+`-Djava.library.path` via `spark.executor.extraJavaOptions`. This **replaces** 
the JVM's
+entire default native library search path. Include every other required native 
library
+directory as well, such as Hadoop native libraries or compression codecs; 
otherwise those
+libraries may become unavailable. Separate directories with the platform path 
separator.
+
+When using `--archives`, Spark extracts the archive to a predictable relative 
path (`./myvenv/`),
+so the path above is stable across executor nodes without any per-node 
configuration.
+
+---
+
+### `spark.executor.cores` and `spark.task.cpus`
+
+Multiple tasks may share an executor, including fractional `spark.task.cpus` 
values.
+Python invocations run one at a time per executor. For throughput, consider 
multiple
+executors with:
+
+```
+spark.executor.cores = 1
+spark.task.cpus      = 1
+```
+
+---
+
+## Supported Types
+
+Common supported Spark SQL input/output types include:
+
+| Category | Types |
+|---|---|
+| Numeric | `ByteType`, `ShortType`, `IntegerType`, `LongType`, `FloatType`, 
`DoubleType` |
+| Boolean | `BooleanType` |
+| String / Binary | `StringType`, `BinaryType` |
+| Temporal | `DateType`, `TimestampType` |
+| Complex | `ArrayType`, `StructType`, `MapType` |
+
+Nested values must satisfy the declared nullability. Map keys cannot be null.
+Only types representable by Spark's Arrow conversion and JVM Arrow accessors 
are
+supported; this is not a guarantee for every Spark SQL type. Unsupported return
+types are rejected on the driver before the function is serialized or tasks 
start.

Review Comment:
   **[Low, docs] This section does not mention that CHAR/VARCHAR return types 
are rejected.**
   
   Since 3b69cd9, `InProcessPythonUDFBuilder.build` rejects CHAR/VARCHAR return 
types with `CHAR_VARCHAR_NOT_SUPPORTED_IN_PYTHON`, as `_check_return_type` does 
on the driver after SPARK-59275, including nested types and UDT storage. The 
table lists only `StringType`, and "Unsupported return types are rejected on 
the driver" does not tell users what to declare instead.
   
   Could this paragraph add one sentence, e.g. "CHAR and VARCHAR return types, 
including nested ones, are rejected; declare STRING instead."?



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalPythonEvaluatorFactory.scala:
##########
@@ -0,0 +1,564 @@
+/*
+ * 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.file.Files
+import java.util.UUID
+import java.util.concurrent.TimeUnit
+import java.util.concurrent.atomic.AtomicInteger
+import java.util.concurrent.locks.ReentrantLock
+
+import scala.collection.mutable.ArrayBuffer
+import scala.jdk.CollectionConverters._
+
+import com.google.common.util.concurrent.Uninterruptibles
+import org.apache.arrow.c.{ArrowArray, ArrowSchema, BaseStruct}
+import org.apache.arrow.util.AutoCloseables
+import org.apache.arrow.vector.VectorSchemaRoot
+
+import org.apache.spark.{SparkEnv, SparkException, TaskContext}
+import org.apache.spark.api.python.ChainedPythonFunctions
+import org.apache.spark.memory.MemoryConsumer
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression, 
JoinedRow, PythonUDF, UnsafeProjection, UnsafeRow}
+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._
+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, unless all of them are UDF arguments that read back from Arrow 
unchanged.
+ * Each batch owns its Arrow buffers so Python can safely retain input arrays.
+ *
+ * The evaluator owns its queue, so that cleanup at task completion is 
coordinated with a
+ * consumer on another thread, such as a pipelined Python writer or a 
TRANSFORM feed thread.
+ */
+class InProcessArrowEvalPythonEvaluatorFactory(
+    childOutput: Seq[Attribute],
+    udfs: Seq[PythonUDF],
+    output: Seq[Attribute],
+    batchSize: Int,
+    maxBytes: Long,
+    timeZoneId: String,
+    largeVarTypes: Boolean,
+    hideTraceback: Boolean,
+    simplifiedTraceback: Boolean,
+    tracebackWithLocals: Boolean,
+    fullValidation: Boolean,
+    metrics: Map[String, SQLMetric])
+  extends EvalPythonEvaluatorFactory(childOutput, udfs, output) {
+
+  private[python] def runtimeSession: 
InProcessPythonRuntime.InterpreterSession =
+    InProcessPythonRuntime.currentSession
+
+  /** Unused: `evaluateJoined` always evaluates the UDFs. */
+  override protected def evaluate(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputSchema: StructType,
+      context: TaskContext): Iterator[InternalRow] =
+    throw SparkException.internalError("In-process UDFs are evaluated with 
their input rows")
+
+  override protected def evaluateJoined(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputs: Seq[Expression],
+      inputSchema: StructType,
+      context: TaskContext): Option[Iterator[InternalRow]] = {
+    import InProcessArrowEvalPythonEvaluatorFactory.{Buffered, ReadBack, 
readsBack}
+    val inputColumns = inputs.length == childOutput.length && 
inputs.zip(childOutput).forall {
+      case (a: Attribute, c) => a.exprId == c.exprId
+      case _ => false
+    }
+    // If all input columns are UDF arguments, they are written to Arrow 
regardless. Read them
+    // back from the exported input vectors instead of buffering every input 
row, if their
+    // values read back from Arrow exactly as written and as fast as an unsafe 
row copy.
+    val joinInput = if (inputColumns && inputSchema.forall(f => 
readsBack(f.dataType))) {
+      ReadBack
+    } else if (inputColumns) {
+      Buffered(None)
+    } else {
+      // Each projected row is written to Arrow before the next input row is 
pulled, so the
+      // arguments go into a reused buffer rather than being copied value by 
value.
+      val projection = UnsafeProjection.create(inputs, childOutput)
+      projection.initialize(context.partitionId())
+      Buffered(Some(projection))
+    }
+    Some(evaluateBatches(funcs, argMetas, rows, inputSchema, context, 
joinInput))
+  }
+
+  private[python] def evaluateBatches(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputSchema: StructType,
+      context: TaskContext,
+      joinInput: InProcessArrowEvalPythonEvaluatorFactory.JoinInput): 
Iterator[InternalRow] = {
+    import InProcessArrowEvalPythonEvaluatorFactory.{Buffered, ReadBack}
+    ArrowUtils.failDuplicatedFieldNames(inputSchema)
+    val functions = funcs.map { case (chain, _) =>
+      if (chain.funcs.size != 1) {
+        throw SparkException.internalError(
+          "In-process UDF chains must use separate evaluation nodes")
+      }
+      chain.funcs.head
+    }
+    val inputOrdinals = argMetas.map(_.map(_.offset))
+    def checkCancellation(): Unit = context.killTaskIfInterrupted()
+
+    val expectedFields = udfs.map { udf =>
+      ArrowUtils.toArrowField("result", udf.dataType, true, timeZoneId, 
largeVarTypes)
+    }
+    val processingTime = new 
InProcessArrowEvalPythonEvaluatorFactory.NanosecondTimer(
+      metrics("pythonProcessingTime"))
+    val initTime = new 
InProcessArrowEvalPythonEvaluatorFactory.NanosecondTimer(
+      metrics("pythonInitTime"))
+    val arrowSchema = ArrowUtils.toArrowSchema(inputSchema, timeZoneId, 
largeVarTypes)
+    // Capture before consuming input: an old task must never join a later 
context's session.
+    val runtime = runtimeSession
+    // Rows are copied out of the queue and Arrow vectors before they are 
returned, so they
+    // remain valid after task completion releases those, on whichever thread 
consumes them.
+    val resultProj = UnsafeProjection.create(output, output)
+    // Spill files go into a directory of the queue's own, created with the 
first disk queue, so
+    // that task completion can delete them when it cannot close the queue.
+    @volatile var spillDir: File = null
+    // Guarded by the queue's monitor.
+    var queueAbandoned = false
+    val (queue, projection) = joinInput match {
+      case Buffered(projection) =>
+        val localDir = new File(Utils.getLocalDir(SparkEnv.get.conf))
+        val serializerManager = SparkEnv.get.serializerManager
+        // Only the consumer holding the iterator's lock adds and removes rows.
+        val queue = new HybridRowQueue(context.taskMemoryManager(), localDir,
+            childOutput.length, serializerManager, lockFree = true) {
+          override protected def createDiskQueue(): RowQueue = synchronized {
+            if (spillDir == null) {
+              spillDir = Files.createTempDirectory(localDir.toPath, 
"inprocess-udf-").toFile
+            }
+            DiskRowQueue(Files.createTempFile(spillDir.toPath, "buffer", 
"").toFile,
+              childOutput.length, serializerManager)
+          }
+
+          // Once task completion leaves the queue to the executor, it must 
not spill for other
+          // consumers into a directory that nothing deletes.
+          override def spill(size: Long, trigger: MemoryConsumer): Long = 
synchronized {
+            if (queueAbandoned) 0L else super.spill(size, trigger)
+          }
+
+          // Queues of a task are distinct memory consumers, whatever their 
case-class fields.
+          override def equals(other: Any): Boolean = this eq 
other.asInstanceOf[AnyRef]
+          override def hashCode(): Int = System.identityHashCode(this)
+          override def canEqual(other: Any): Boolean = false
+        }
+        (queue, projection.orNull)
+      case ReadBack => (null, null)
+    }
+    val joined = new JoinedRow
+    val handles = functions.map(_ => UUID.randomUUID().toString)
+    var registered = false
+    var writer: ArrowWriter = null
+    val results = ArrayBuffer.empty[ArrowColumnVector]
+    var startedAt = 0L
+
+    def closeBatch(): Unit = {
+      val resources = ArrayBuffer.empty[AutoCloseable]
+      resources ++= results
+      results.clear()
+      if (writer != null) {
+        resources += writer.root
+        writer = null
+      }
+      AutoCloseables.close(resources.asJava)
+    }
+
+    val resources = new 
InProcessArrowEvalPythonEvaluatorFactory.IteratorResources(
+      // Closing the queue deletes the spill files it tracks; deleteQuietly 
also removes any
+      // other, without starting a process or throwing, also on an interrupted 
thread.
+      releaseTaskMemory = () => if (queue != null) {
+        Utils.tryWithSafeFinally(queue.close())(Utils.deleteQuietly(spillDir))
+      },
+      abandonTaskMemory = () => if (queue != null) {
+        queue.synchronized {
+          queueAbandoned = true
+          Utils.deleteQuietly(spillDir)
+        }
+      },
+      releaseOthers = () => {
+        if (startedAt != 0L) {
+          metrics("pythonTotalTime") += (System.nanoTime() - startedAt) / 
1000000
+        }
+        Utils.tryWithSafeFinally {
+          closeBatch()
+        } {
+          if (registered) runtime.release(handles)
+        }
+      })
+
+    context.addTaskCompletionListener[Unit](_ => resources.close())
+
+    new Iterator[InternalRow] {
+      private var batchIter: Iterator[InternalRow] = Iterator.empty
+
+      private def endOfInput: Nothing =
+        throw new NoSuchElementException("End of in-process UDF input")
+
+      // Releases the resources on failure without replacing its exception.
+      private def fail(t: Throwable): Nothing =
+        Utils.tryWithSafeFinally { throw t } { resources.close() }
+
+      // Called with the lock held.
+      private def hasNextLocked: Boolean = {
+        if (startedAt == 0L) startedAt = System.nanoTime()
+        checkCancellation()
+        val available = batchIter.hasNext || {
+          var open = false
+          val more = try {
+            resources.startReadingInput() && rows.hasNext
+          } finally {
+            open = resources.endReadingInput()
+          }
+          more && open
+        }
+        if (!available) resources.close()
+        available
+      }
+
+      // Each call takes the lock without allocating a closure per row. Within 
a batch, only
+      // the consumer advances `batchIter`, whose `hasNext` compares row 
indexes.
+      override def hasNext: Boolean = {
+        if (batchIter.hasNext && !resources.isClosed) return true
+        if (!resources.enter()) return false
+        val available = try {
+          try hasNextLocked catch { case t: Throwable => fail(t) }
+        } finally {
+          resources.exit()
+        }
+        // Task completion may have closed the iterator while the input was 
being read.
+        available && !resources.isClosed
+      }
+
+      override def next(): InternalRow = {
+        if (!resources.enter()) endOfInput
+        try {
+          try {
+            if (!hasNextLocked) endOfInput
+            if (!batchIter.hasNext) nextBatch()
+            val result = batchIter.next()
+            resultProj(if (queue != null) joined(queue.remove(), result) else 
result)
+          } catch {
+            case t: Throwable => fail(t)
+          }
+        } finally {
+          resources.exit()
+        }
+      }
+
+      // Runs Python without the lock unless task completion already happened, 
and ends the
+      // input instead of returning the result if it happens meanwhile.
+      private def python[T](body: => T): T = {
+        if (resources.isClosed) endOfInput
+        val result = resources.withoutLock(body)
+        if (resources.isClosed) endOfInput
+        result
+      }
+
+      /**
+       * Writes the next input row to the batch, returning false at the end of 
input or once
+       * task completion happened. If it happens while the row is read, the 
row is dropped.
+       */
+      private def pullRow(): Boolean = {
+        var open = false
+        val row = try {
+          // Checked again after `hasNext`, which may wait for input, not to 
read another row.
+          if (resources.startReadingInput() && rows.hasNext && 
!resources.isClosed) {
+            rows.next()
+          } else {
+            null
+          }
+        } finally {
+          open = resources.endReadingInput()
+        }
+        if (!open) endOfInput
+        if (row == null) return false
+        if (queue != null) queue.add(row.asInstanceOf[UnsafeRow])
+        writer.write(if (projection != null) projection(row) else row)
+        true
+      }
+
+      // Called with the lock held.
+      private def nextBatch(): Unit = {
+        closeBatch()
+        val root = VectorSchemaRoot.create(arrowSchema, 
ArrowUtils.rootAllocator)
+        writer = try {
+          ArrowWriter.create(root)
+        } catch {
+          case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
root.close() }
+        }
+        // Task completion stops the fill within a row, and Python never sees 
a partial batch.
+        var count = 0
+        while ((batchSize <= 0 || count < batchSize) &&
+            (count == 0 || writer.sizeInBytes() < maxBytes) && {
+              checkCancellation()
+              pullRow()
+            }) {
+          count += 1
+        }
+        if (resources.isClosed) endOfInput
+        if (!registered) {
+          // Mark before registering so failure after any registration still 
cleans up.
+          registered = true
+          functions.indices.foreach { i =>
+            val func = functions(i)
+            initTime.add(python(runtime.register(handles(i), 
func.command.toArray,
+              expectedFields(i), func.pythonVer, hideTraceback, 
simplifiedTraceback,
+              tracebackWithLocals, fullValidation)))
+          }
+        }
+        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 track[S <: BaseStruct](struct: S): S = {
+            val closer: AutoCloseable = () => 
InProcessArrowBridge.closeStruct(struct)
+            structs += closer
+            struct
+          }
+          def array(): ArrowArray = 
track(ArrowArray.allocateNew(ArrowUtils.rootAllocator))
+          def schema(): ArrowSchema = 
track(ArrowSchema.allocateNew(ArrowUtils.rootAllocator))
+          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))
+            }
+            processingTime.add(python(runtime.invoke(

Review Comment:
   **[Low] Python time is lost for a batch whose Python returns after task 
completion.**
   
   When task completion happens during the call, 
`processingTime.add(python(runtime.invoke(...)))` adds nothing: `python()` 
throws `endOfInput` when Python returns, before `add` runs. `pythonTotalTime` 
(L211) has a similar gap. When the listener finds the consumer in Python, the 
consumer adds it in `releaseOthers` after Python returns, and by then the task 
may have already reported its accumulator updates, so the addition is never 
reported.
   
   This happens only with a consumer on another thread (a pipelined writer or a 
TRANSFORM feed thread) that is in Python at completion, e.g. under a LIMIT, but 
those are exactly the slow batches. Suggestion: add the elapsed time before the 
`isClosed` check in `python()`, e.g. by letting it take the timer, and record 
`pythonTotalTime` when the listener requests the close rather than when the 
last resource is released.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,419 @@
+/*
+ * 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.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+import java.util.concurrent.atomic.AtomicInteger
+
+import scala.collection.mutable
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, 
NamingConventionClassEnquirer, PyConfig, SharedInterpreter}
+import org.apache.arrow.c.{ArrowSchema, Data}
+import org.apache.arrow.vector.types.pojo.Field
+
+import org.apache.spark.{TaskContext, TaskKilledException}
+import org.apache.spark.api.python.{PythonException, PythonUtils}
+import org.apache.spark.internal.Logging
+import org.apache.spark.internal.config.Python
+import org.apache.spark.sql.util.ArrowUtils
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var mainConfigured = false
+  @volatile private var sharedConfigured = false
+  @volatile private var bootstrappedSitePackages: Option[Seq[String]] = None
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  // Keep JEP references out of the singleton's verifier so currentSession can 
report
+  // an uninitialized runtime even when the provided JEP JAR is absent.
+  private[python] object InterpreterConfiguration {
+    def configure(sitePackages: Seq[String]): Unit = {
+      if (!mainConfigured) {
+        // 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(
+          
PyConfig.isolated().setUseEnvironment(false).setHashSeed(0).setUseHashSeed(true))
+        mainConfigured = true
+      }
+      if (!sharedConfigured) {
+        // JEP imports its Python package during construction, before our 
bootstrap runs.
+        SharedInterpreter.setConfig(interpreterConfig(sitePackages))
+      }
+    }
+
+    def interpreterConfig(sitePackages: Seq[String]): JepConfig = {
+      require(sitePackages.forall(Python.isValidInProcessPath), 
Python.IN_PROCESS_PATH_RULE)
+      val config = new JepConfig().setClassEnquirer(new 
NamingConventionClassEnquirer(false))
+      // Calling addIncludePaths with no arguments adds the working directory 
in JEP.
+      if (sitePackages.nonEmpty) config.addIncludePaths(sitePackages: _*)
+      config
+    }
+  }
+
+  private class ManagedSharedInterpreter extends SharedInterpreter {
+    override protected def configureInterpreter(config: JepConfig): Unit = {
+      // JEP invokes this hook after native initialization, from its 
constructor. Close
+      // here if configuration fails, before the caller can receive an 
interpreter handle.
+      try {
+        super.configureInterpreter(config)
+        sharedConfigured = true
+      } catch {
+        case t: Throwable => Utils.tryWithSafeFinally { throw t } { close() }
+      }
+    }
+  }
+
+  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: ' + " +
+      "ascii(type(_bootstrap_error).__name__ + ': ' + str(_bootstrap_error))) 
from None\n"
+  }
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    bootstrappedSitePackages.foreach { paths =>
+      if (paths != sitePackages) {
+        throw new LifecycleException("In-process Python has already configured 
different " +
+          "sitePackages. Restart the executor process before changing 
interpreter configuration.")
+      }
+    }
+    if (active != null && !active.isTerminated) {
+      active.requireCompatible(sitePackages)
+    } else {
+      InterpreterConfiguration.configure(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)
+    // `shutdown` keeps the stopped session, so a session that is not running 
was stopped.
+    checkState(active.isRunning, StoppedMessage)
+    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 val StoppedMessage =
+    "In-process Python has been stopped (executor or SparkContext shutdown)"
+
+  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
+    // Calls submitted to the interpreter thread that have not finished or 
been cancelled.
+    private val pendingCalls = new AtomicInteger()
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+    // Guarded by this session's monitor. Shutdown must keep Python-owned 
result buffers
+    // pinned until their tasks have released the JVM CDI references.
+    private val registeredHandles = mutable.Set.empty[String]
+
+    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) {

Review Comment:
   **[Low, cleanup] This branch cannot be reached from `initialize`.**
   
   A session becomes `active` only after `InterpreterSession.initialize()` has 
set `bootstrappedSitePackages = Some(sitePackages)` (L234), and nothing resets 
`bootstrappedSitePackages`. `InProcessPythonRuntime.initialize` (L95-100) 
throws on any difference from it before it calls `requireCompatible`, so 
`sitePackages == paths` always holds here for the `active` session. Only 
`InProcessPythonRuntimeSuite` ("lifecycle errors distinguish configuration 
mismatch from stopping") reaches this branch, by calling `requireCompatible` 
directly.
   
   Could `requireCompatible` keep only the stopping check, so that a single 
check and message covers the configuration mismatch?



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