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


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalExec.scala:
##########
@@ -0,0 +1,198 @@
+/*
+ * 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 scala.collection.mutable.ArrayBuffer
+import scala.jdk.CollectionConverters._
+
+import org.apache.arrow.c.{ArrowArray, ArrowSchema}
+import org.apache.arrow.util.AutoCloseables
+import org.apache.arrow.vector.VectorSchemaRoot
+
+import org.apache.spark.TaskContext
+import org.apache.spark.rdd.RDD
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression, 
UnsafeProjection}
+import org.apache.spark.sql.execution.{SparkPlan, UnaryExecNode}
+import org.apache.spark.sql.execution.arrow.ArrowWriter
+import org.apache.spark.sql.types.{StructField, StructType}
+import org.apache.spark.sql.util.ArrowUtils
+import org.apache.spark.sql.vectorized.{ArrowColumnVector, ColumnarBatch, 
ColumnVector}
+import org.apache.spark.util.Utils
+
+/**
+ * Evaluates scalar Python UDFs using Arrow CDI in the executor process. Rows 
(including
+ * computed UDF arguments) are written to Arrow once, then input and output 
buffers cross
+ * the JVM/Python boundary without IPC serialization. The runtime owns the JEP 
thread.
+ */
+case class InProcessArrowEvalExec(
+    udfs: Seq[InProcessPythonUDF],
+    resultAttrs: Seq[Attribute],
+    child: SparkPlan) extends UnaryExecNode {
+
+  override def output: Seq[Attribute] = child.output ++ resultAttrs
+
+  override protected def doExecute(): RDD[InternalRow] = {
+    val expressions = ArrayBuffer[Expression](child.output: _*)
+    val inputOrdinals = udfs.map { udf =>
+      udf.children.map {
+        case attr: Attribute if child.output.exists(_.exprId == attr.exprId) =>
+          child.output.indexWhere(_.exprId == attr.exprId)
+        case expr =>
+          expressions += expr
+          expressions.size - 1
+      }
+    }
+    // Synthetic names also allow joins with duplicate output column names.
+    val inputSchema = StructType(expressions.zipWithIndex.map { case (expr, i) 
=>
+      StructField(s"_input$i", expr.dataType, expr.nullable)
+    }.toSeq)
+    val inputExpressions = expressions.toSeq
+    val childOutput = child.output
+    val resultOutput = output
+    val batchSize = conf.arrowMaxRecordsPerBatch
+    val maxBytes = conf.arrowMaxBytesPerBatch
+    val timeZoneId = conf.sessionLocalTimeZone
+
+    child.execute().mapPartitions { rows =>
+      val context = Option(TaskContext.get())
+      def checkCancellation(): Unit = 
context.foreach(_.killTaskIfInterrupted())
+
+      val resultProjection = UnsafeProjection.create(resultOutput, 
resultOutput)
+      val projectInput: InternalRow => InternalRow =
+        if (inputExpressions.size == childOutput.size) {
+          identity[InternalRow]
+        } else {
+          val projection = UnsafeProjection.create(inputExpressions, 
childOutput)
+          projection.initialize(TaskContext.getPartitionId())
+          projection
+        }
+      val root = VectorSchemaRoot.create(
+        ArrowUtils.toArrowSchema(inputSchema, timeZoneId, false), 
ArrowUtils.rootAllocator)
+      val writer = try {
+        ArrowWriter.create(root)
+      } catch {
+        case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
root.close() }
+      }
+      val results = ArrayBuffer.empty[ArrowColumnVector]
+      var closed = false
+
+      def closeResults(): Unit = {
+        val previous = results.toArray
+        results.clear()
+        AutoCloseables.close(previous: _*)
+      }
+
+      def close(): Unit = {
+        if (!closed) {
+          closed = true
+          Utils.tryWithSafeFinally { closeResults() } { writer.root.close() }
+        }
+      }
+
+      context.foreach(_.addTaskCompletionListener[Unit](_ => close()))
+
+      new Iterator[InternalRow] {
+        private var batchIter: Iterator[InternalRow] = Iterator.empty
+
+        override def hasNext: Boolean = {
+          checkCancellation()
+          val available = !closed && (batchIter.hasNext || rows.hasNext)
+          if (!available) close()
+          available
+        }
+
+        override def next(): InternalRow = {
+          if (!hasNext) throw new NoSuchElementException("End of in-process 
UDF input")
+          try {
+            if (!batchIter.hasNext) {
+              closeResults()
+              writer.reset()

Review Comment:
   Each batch now gets a fresh vector root and writer; buffers exported to 
Python are no longer reset and reused for later batches. Retained Python arrays 
keep their buffers alive through the CDI ownership callbacks. Added a stateful 
regression that retains prior inputs and checks their values and nulls after 
subsequent batches. The docs also explain that retaining arrays retains native 
memory.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonChecks.scala:
##########
@@ -0,0 +1,63 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
+import org.apache.spark.sql.catalyst.rules.Rule
+import org.apache.spark.sql.internal.SQLConf
+
+/**
+ * Validates that in-process Python UDFs are only used when exactly one task 
can run per
+ * executor, preventing GIL contention on [[InProcessPythonRuntime]]'s shared 
interpreter.
+ *
+ * The constraint: spark.executor.cores / spark.task.cpus == 1
+ *
+ * Typical correct configuration:
+ *   spark.executor.cores=1  (one core per executor, parallelism via more 
executors)
+ *
+ * Runs after [[ExtractInProcessPythonUDFs]] in the "Extract InProcess Python 
UDFs" optimizer
+ * batch, so it sees [[InProcessEvalPython]] nodes.
+ */
+object InProcessPythonChecks extends Rule[LogicalPlan] {
+
+  override def apply(plan: LogicalPlan): LogicalPlan = {
+    plan.foreach {
+      case _: InProcessEvalPython => checkConcurrencyConfig()
+      case _ =>
+    }
+    plan
+  }
+
+  private def checkConcurrencyConfig(): Unit = {
+    val conf = SQLConf.get
+    val executorCores =
+      conf.getConfString("spark.executor.cores", "1").toInt
+    val taskCpus =
+      conf.getConfString("spark.task.cpus", "1").toInt

Review Comment:
   Removed the planning-time CPU check. Calls still run on the dedicated 
interpreter thread and are serialized by the existing lock, so multiple tasks 
can share an executor. Function instances are now scoped to each task, while 
imported modules can still have executor-wide state. The integration suite 
passes with `local[2]` and `spark.task.cpus=0.5`; the docs describe one task 
per executor as a throughput recommendation rather than a correctness 
requirement.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonUDF.scala:
##########
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import org.apache.spark.sql.catalyst.expressions.{
+  Attribute, AttributeReference, AttributeSet, Expression, ExprId, 
NamedExpression, Unevaluable
+}
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, UnaryNode}
+import org.apache.spark.sql.types.DataType
+
+/**
+ * A cloudpickle-serialized Python function to be executed in-process via jep
+ * (Java Embedded Python).
+ *
+ * Distinct from [[org.apache.spark.sql.catalyst.expressions.PythonUDF]] which 
uses an
+ * out-of-process Python worker connected via socket.
+ *
+ * Evaluated by [[InProcessArrowEvalExec]], which passes Arrow column buffers 
to CPython
+ * as PyArrow arrays via native memory addresses (zero-copy input), then 
imports the
+ * PyArrow result buffers through CDI without copying.
+ *
+ * @param name            display name for plan explain output
+ * @param serializedFunc  cloudpickle-serialized Python function bytes
+ * @param children        input column expressions
+ * @param dataType        declared return type (validated against the Arrow 
result)
+ * @param udfDeterministic whether the UDF is deterministic
+ * @param resultId        unique identifier for this UDF result
+ */
+case class InProcessPythonUDF(
+    name: String,
+    serializedFunc: Array[Byte],
+    children: Seq[Expression],
+    dataType: DataType,
+    udfDeterministic: Boolean = true,
+    resultId: ExprId = NamedExpression.newExprId)
+  extends Expression with Unevaluable {
+
+  override def nullable: Boolean = true
+  override def prettyName: String = name
+
+  override lazy val deterministic: Boolean =
+    udfDeterministic && children.forall(_.deterministic)
+
+  lazy val resultAttribute: Attribute =
+    AttributeReference(name, dataType, nullable)(exprId = resultId)
+
+  override def toString: String = s"$name(${children.mkString(", 
")})#${resultId.id}"
+
+  override protected def withNewChildrenInternal(
+      newChildren: IndexedSeq[Expression]): InProcessPythonUDF =
+    copy(children = newChildren)
+}
+
+/**
+ * Logical plan node that evaluates [[InProcessPythonUDF]]s in-process via jep.
+ * Inserted by [[ExtractInProcessPythonUDFs]] during query optimization, 
before physical planning.
+ * Planned as [[InProcessArrowEvalExec]] by 
[[org.apache.spark.sql.execution.SparkStrategies]].
+ */
+case class InProcessEvalPython(

Review Comment:
   Replaced `InProcessEvalPython` with `ArrowEvalPython` carrying the 
in-process eval type. Existing predicate/limit pushdown and scan rules can now 
recognize the node. Added plan tests for ordinary filter and limit pushdown, 
plus an integration test where the non-UDF predicate removes zero denominators 
before the division UDF runs.



##########
python/pyspark/inprocess/udf.py:
##########
@@ -0,0 +1,175 @@
+#
+# 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()
+"""
+
+from typing import Callable
+
+import pyarrow as pa
+
+from pyspark import cloudpickle
+from pyspark.sql.types import (
+    BooleanType,
+    ByteType,
+    DataType,
+    DoubleType,
+    FloatType,
+    IntegerType,
+    LongType,
+    ShortType,
+)
+
+# Map from Spark SQL DataType to PyArrow type for output type enforcement.
+_SPARK_TO_ARROW: dict = {
+    LongType(): pa.int64(),
+    IntegerType(): pa.int32(),
+    DoubleType(): pa.float64(),
+    FloatType(): pa.float32(),
+    BooleanType(): pa.bool_(),
+    ShortType(): pa.int16(),
+    ByteType(): pa.int8(),
+}
+
+
+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 ``InProcessPythonUDF``
+    on the JVM side.
+    """
+
+    def __init__(self, func: Callable, return_type: DataType, deterministic: 
bool = True) -> None:
+        self._return_type: DataType = return_type
+        self._deterministic: bool = deterministic
+        self._name: str = getattr(func, "__name__", "inprocess_udf")
+
+        # Wrap the function to cast its output to the declared return type.
+        # This handles the case where the UDF's input column type differs from
+        # the declared return type (e.g. input is int64, return_type is 
IntegerType).
+        arrow_type = _SPARK_TO_ARROW.get(return_type)
+        if arrow_type is not None:
+
+            def _wrapped(*args, _fn=func, _atype=arrow_type):
+                result = _fn(*args)
+                if not isinstance(result, pa.Array):
+                    raise TypeError("In-process UDF must return a 
pyarrow.Array")
+                if result.type != _atype:
+                    result = result.cast(_atype)
+                return result
+
+            self._serialized: bytes = cloudpickle.dumps(_wrapped)

Review Comment:
   Added a driver/embedded Python major.minor version check before unpickling, 
using `PYTHON_VERSION_MISMATCH`. Captured Spark `Broadcast` and `Accumulator` 
objects are rejected during serialization, with regression tests. Broadcasts, 
accumulators, `addPyFile`, and Python `TaskContext` are documented as 
unsupported; the docs explicitly include access through imported modules, which 
the closure pickler cannot reliably detect. Dependencies must be installed on 
executors before startup, optionally using 
`spark.inprocess.python.sitePackages`.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonUDF.scala:
##########
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import org.apache.spark.sql.catalyst.expressions.{
+  Attribute, AttributeReference, AttributeSet, Expression, ExprId, 
NamedExpression, Unevaluable
+}
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, UnaryNode}
+import org.apache.spark.sql.types.DataType
+
+/**
+ * A cloudpickle-serialized Python function to be executed in-process via jep
+ * (Java Embedded Python).
+ *
+ * Distinct from [[org.apache.spark.sql.catalyst.expressions.PythonUDF]] which 
uses an
+ * out-of-process Python worker connected via socket.
+ *
+ * Evaluated by [[InProcessArrowEvalExec]], which passes Arrow column buffers 
to CPython
+ * as PyArrow arrays via native memory addresses (zero-copy input), then 
imports the
+ * PyArrow result buffers through CDI without copying.
+ *
+ * @param name            display name for plan explain output
+ * @param serializedFunc  cloudpickle-serialized Python function bytes
+ * @param children        input column expressions
+ * @param dataType        declared return type (validated against the Arrow 
result)
+ * @param udfDeterministic whether the UDF is deterministic
+ * @param resultId        unique identifier for this UDF result
+ */
+case class InProcessPythonUDF(
+    name: String,
+    serializedFunc: Array[Byte],
+    children: Seq[Expression],
+    dataType: DataType,
+    udfDeterministic: Boolean = true,
+    resultId: ExprId = NamedExpression.newExprId)
+  extends Expression with Unevaluable {
+
+  override def nullable: Boolean = true
+  override def prettyName: String = name
+
+  override lazy val deterministic: Boolean =

Review Comment:
   Using `PythonUDF` now provides the existing `UserDefinedExpression` contract 
and propagates the deterministic flag, so `PullOutNondeterministic` handles 
these expressions. Added plan and integration tests for nondeterministic 
in-process UDFs in grouping keys and sort expressions.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonUDF.scala:
##########
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import org.apache.spark.sql.catalyst.expressions.{
+  Attribute, AttributeReference, AttributeSet, Expression, ExprId, 
NamedExpression, Unevaluable
+}
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, UnaryNode}
+import org.apache.spark.sql.types.DataType
+
+/**
+ * A cloudpickle-serialized Python function to be executed in-process via jep
+ * (Java Embedded Python).
+ *
+ * Distinct from [[org.apache.spark.sql.catalyst.expressions.PythonUDF]] which 
uses an
+ * out-of-process Python worker connected via socket.
+ *
+ * Evaluated by [[InProcessArrowEvalExec]], which passes Arrow column buffers 
to CPython
+ * as PyArrow arrays via native memory addresses (zero-copy input), then 
imports the
+ * PyArrow result buffers through CDI without copying.
+ *
+ * @param name            display name for plan explain output
+ * @param serializedFunc  cloudpickle-serialized Python function bytes
+ * @param children        input column expressions
+ * @param dataType        declared return type (validated against the Arrow 
result)
+ * @param udfDeterministic whether the UDF is deterministic
+ * @param resultId        unique identifier for this UDF result
+ */
+case class InProcessPythonUDF(
+    name: String,
+    serializedFunc: Array[Byte],

Review Comment:
   The builder now creates `PythonUDF` with a `SimplePythonFunction` whose 
command bytes use content equality. This also reuses `PythonUDF`'s result-ID 
canonicalization and `PythonFuncExpression.expensive`. Added tests for repeated 
calls in grouping expressions, semantic equality across rebuilt queries, and 
sharing deterministic duplicate calls.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,96 @@
+#
+# 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.
+#
+
+
+"""
+In-process Python UDF runtime entry point.
+
+``_inprocess_invoke`` is imported into the jep SharedInterpreter's global 
namespace
+during executor initialization (see ``InProcessPythonRuntime.initialize()``), 
then called
+directly from the JVM via ``interp.invoke("_inprocess_invoke", ...)``.
+
+Both input and output use the Arrow C Data Interface (CDI). The JVM 
pre-allocates
+ArrowArray/ArrowSchema C structs for every input column and for the output, 
passing
+their native addresses as Python ints. Input arrays are reconstructed via
+``pa.Array._import_from_c`` (zero-copy). The output is written via 
``arr._export_to_c``
+into the JVM-owned structs (zero-copy).
+
+jep type conversions (Java -> Python):
+    byte[]                -> bytes (or sequence of signed ints; masked to 
unsigned below)
+    List<Long> (boxed)    -> list of Python ints
+    Long                  -> int
+"""
+
+import traceback as _traceback
+from functools import lru_cache
+
+import pyarrow as pa
+
+from pyspark import cloudpickle
+from pyspark.sql.pandas.types import to_arrow_type
+from pyspark.sql.types import _parse_datatype_json_string
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+
+
+@lru_cache(maxsize=128)
+def _load_udf(serialized_udf: bytes, return_type_json: str, timezone: str):
+    return (
+        cloudpickle.loads(serialized_udf),
+        to_arrow_type(_parse_datatype_json_string(return_type_json), 
timezone=timezone),
+    )
+
+
+def _validate_result(result, expected_rows: int, expected_type: pa.DataType) 
-> None:
+    if not isinstance(result, pa.Array):
+        raise TypeError(f"In-process UDF must return a pyarrow.Array, got 
{type(result).__name__}")
+    if len(result) != expected_rows:
+        raise ValueError(f"In-process UDF returned {len(result)} rows; 
expected {expected_rows}")
+    if result.type != expected_type:

Review Comment:
   The runtime now compares structural types with nested nullability 
normalized, checks actual values against the declared non-nullable fields, and 
casts compatible results to the declared Arrow type. Null children beneath null 
parent structs are ignored. Added coverage for list, struct, and map 
nullability, including rejection of actual nulls in required fields.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,96 @@
+#
+# 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.
+#
+
+
+"""
+In-process Python UDF runtime entry point.
+
+``_inprocess_invoke`` is imported into the jep SharedInterpreter's global 
namespace
+during executor initialization (see ``InProcessPythonRuntime.initialize()``), 
then called
+directly from the JVM via ``interp.invoke("_inprocess_invoke", ...)``.
+
+Both input and output use the Arrow C Data Interface (CDI). The JVM 
pre-allocates
+ArrowArray/ArrowSchema C structs for every input column and for the output, 
passing
+their native addresses as Python ints. Input arrays are reconstructed via
+``pa.Array._import_from_c`` (zero-copy). The output is written via 
``arr._export_to_c``
+into the JVM-owned structs (zero-copy).
+
+jep type conversions (Java -> Python):
+    byte[]                -> bytes (or sequence of signed ints; masked to 
unsigned below)
+    List<Long> (boxed)    -> list of Python ints
+    Long                  -> int
+"""
+
+import traceback as _traceback
+from functools import lru_cache
+
+import pyarrow as pa
+
+from pyspark import cloudpickle
+from pyspark.sql.pandas.types import to_arrow_type
+from pyspark.sql.types import _parse_datatype_json_string
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+
+
+@lru_cache(maxsize=128)
+def _load_udf(serialized_udf: bytes, return_type_json: str, timezone: str):
+    return (
+        cloudpickle.loads(serialized_udf),
+        to_arrow_type(_parse_datatype_json_string(return_type_json), 
timezone=timezone),
+    )
+
+
+def _validate_result(result, expected_rows: int, expected_type: pa.DataType) 
-> None:
+    if not isinstance(result, pa.Array):
+        raise TypeError(f"In-process UDF must return a pyarrow.Array, got 
{type(result).__name__}")
+    if len(result) != expected_rows:
+        raise ValueError(f"In-process UDF returned {len(result)} rows; 
expected {expected_rows}")
+    if result.type != expected_type:
+        raise TypeError(f"In-process UDF returned {result.type}; expected 
{expected_type}")
+    result.validate()
+
+
+def _inprocess_invoke(
+    serialized_udf,
+    input_array_ptrs,
+    input_schema_ptrs,
+    output_array_ptr: int,
+    output_schema_ptr: int,
+    expected_rows: int,
+    return_type_json: str,
+    timezone: str,
+) -> None:
+    """Consume input CDI structs and export a validated, row-preserving result.
+
+    The caller owns the struct memory and releases any unconsumed exports on 
failure.
+    Imported input arrays and exported output buffers follow Arrow's release 
callbacks.
+    """
+    udf_key = bytes(b & 0xFF for b in serialized_udf)

Review Comment:
   Replaced the per-batch serialized-closure lookup with task-scoped 
registration. The command is converted and unpickled once per task/UDF, and 
subsequent batches pass only a small handle and CDI addresses. Return-type JSON 
is computed outside the partition loop, and its Arrow type is parsed once 
during registration. Task completion releases the handle, including 
cancellation cleanup. This also gives each task its own function state instead 
of sharing a cached closure across tasks. Added registration/state-isolation 
tests; I haven't rerun the large-closure benchmark yet.



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