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


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessArrowEvalPythonEvaluatorFactory.scala:
##########
@@ -0,0 +1,237 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.util.UUID
+
+import scala.collection.mutable.ArrayBuffer
+import scala.jdk.CollectionConverters._
+
+import org.apache.arrow.c.{ArrowArray, ArrowSchema}
+import org.apache.arrow.util.AutoCloseables
+import org.apache.arrow.vector.VectorSchemaRoot
+
+import org.apache.spark.{SparkException, TaskContext}
+import org.apache.spark.api.python.ChainedPythonFunctions
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{Attribute, PythonUDF}
+import org.apache.spark.sql.execution.arrow.ArrowWriter
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.execution.python.EvalPythonExec.ArgumentMetadata
+import org.apache.spark.sql.types.StructType
+import org.apache.spark.sql.util.ArrowUtils
+import org.apache.spark.sql.vectorized.{ArrowColumnVector, ColumnarBatch, 
ColumnVector}
+import org.apache.spark.util.Utils
+
+/**
+ * Evaluates scalar Python UDFs using Arrow CDI in the executor process. Only 
UDF arguments
+ * are converted to Arrow. Original rows are buffered in a spillable queue and 
joined with
+ * the results. Each batch owns its Arrow buffers so Python can safely retain 
input arrays.
+ */
+class InProcessArrowEvalPythonEvaluatorFactory(
+    childOutput: Seq[Attribute],
+    udfs: Seq[PythonUDF],
+    output: Seq[Attribute],
+    batchSize: Int,
+    maxBytes: Long,
+    timeZoneId: String,
+    largeVarTypes: Boolean,
+    hideTraceback: Boolean,
+    simplifiedTraceback: Boolean,
+    tracebackWithLocals: Boolean,
+    metrics: Map[String, SQLMetric])
+  extends EvalPythonEvaluatorFactory(childOutput, udfs, output) {
+
+  override protected def evaluate(
+      funcs: Seq[(ChainedPythonFunctions, Long)],
+      argMetas: Array[Array[ArgumentMetadata]],
+      rows: Iterator[InternalRow],
+      inputSchema: StructType,
+      context: TaskContext): Iterator[InternalRow] = {
+    ArrowUtils.failDuplicatedFieldNames(inputSchema)
+    val functions = funcs.map { case (chain, _) =>
+      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)
+    var runtime: InProcessPythonRuntime.InterpreterSession = null
+    val handles = functions.map(_ => UUID.randomUUID().toString)
+    var registered = false
+    var writer: ArrowWriter = null
+    val results = ArrayBuffer.empty[ArrowColumnVector]
+    var closed = false
+    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)

Review Comment:
   The runtime now retains each handle's exported result until the next 
invocation or task release, after the JVM has closed its CDI references. Those 
Python references are dropped on the interpreter thread. Shutdown also tracks 
outstanding registrations so it cannot clear retained results before task 
cleanup.
   
   Added pure-Python CDI and Spark/NumPy finalizer tests that verify the 
cleanup thread, plus coverage for bounded shutdown followed by late task 
cleanup.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,344 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.io.File
+import java.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, PyConfig, 
SharedInterpreter}
+import org.apache.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.sql.util.ArrowUtils
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  val SITE_PACKAGES_CONFIG = "spark.inprocess.python.sitePackages"
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var mainConfigured = false
+  @volatile private var sharedConfigured = false
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  private def configureInterpreter(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(new 
JepConfig().addIncludePaths(sitePackages: _*))

Review Comment:
   Configured `NamingConventionClassEnquirer(false)` to avoid `ClassList` 
scanning and Java package shadowing. `addIncludePaths` is now called only for 
nonempty paths. Validation rejects single quotes, backslashes, newlines, and 
the platform path separator before passing paths to JEP. Added Python-package 
import and path-validation coverage.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,310 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+
+"""Arrow CDI entry points called on the executor's dedicated JEP interpreter 
thread.
+
+Functions are registered once per task and released when that task finishes. 
Calls
+pass only a handle and CDI addresses, so large closures are not copied per 
batch.
+"""
+
+import sys
+from typing import Any, Callable, Iterable, Optional, Sequence
+
+import pyarrow as pa
+import pyarrow.compute as pc
+
+from pyspark import cloudpickle
+from pyspark.errors import PySparkRuntimeError
+from pyspark.sql.pandas.utils import require_minimum_pyarrow_version
+from pyspark.util import _format_exception
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+NullChecker = Callable[[pa.Array], None]
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType, NullChecker, 
bool, bool, bool]] = {}
+
+
+def _jep_safe_message(message: str) -> str:
+    # JEP uses JNI modified UTF-8 for exception text. Keep the transport ASCII 
and
+    # escape NUL explicitly; ordinary UTF-8 and embedded NUL are not safe here.
+    return message.encode("ascii", 
"backslashreplace").decode("ascii").replace("\0", "\\x00")
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    schema_ptr: int,
+    python_version: str,
+    hide_traceback: bool = False,
+    simplified_traceback: bool = False,
+    traceback_with_locals: bool = False,
+) -> None:
+    try:
+        require_minimum_pyarrow_version()
+        embedded_version = "%d.%d" % sys.version_info[:2]
+        if python_version != embedded_version:
+            raise PySparkRuntimeError(
+                errorClass="PYTHON_VERSION_MISMATCH",
+                messageParameters={
+                    "worker_version": embedded_version,
+                    "driver_version": python_version,
+                },
+            )
+        # JEP exposes direct ByteBuffers through the buffer protocol. Unpickle 
a separate
+        # function per task without iterating over a PyJArray one JNI call per 
byte.
+        func = cloudpickle.loads(memoryview(serialized_udf))
+        # The JVM is the single source of truth for Arrow layout and logical 
metadata.
+        expected_type = pa.Field._import_from_c(schema_ptr).type
+        checker = _null_checker(expected_type) or (lambda array: None)
+        _udfs[handle] = (
+            func,
+            expected_type,
+            checker,
+            hide_traceback,
+            simplified_traceback,
+            traceback_with_locals,
+        )
+    except BaseException as error:
+        # In JEP, an uncaught SystemExit can terminate the entire executor JVM.
+        raise RuntimeError(
+            _UDF_TRACEBACK_SENTINEL
+            + _jep_safe_message(
+                _format_exception(
+                    error, hide_traceback, simplified_traceback, 
traceback_with_locals
+                )
+            )
+        ) from None
+
+
+def _inprocess_release(handles: Iterable[str]) -> None:
+    for handle in handles:
+        _udfs.pop(handle, None)
+
+
+def _nullable_type(data_type: pa.DataType) -> pa.DataType:
+    def nullable_field(field: pa.Field) -> pa.Field:
+        return pa.field(field.name, _nullable_type(field.type), nullable=True)
+
+    if pa.types.is_struct(data_type):
+        return pa.struct([nullable_field(field) for field in data_type])
+    if pa.types.is_list(data_type):
+        return pa.list_(nullable_field(data_type.value_field))
+    if pa.types.is_large_list(data_type):
+        return pa.large_list(nullable_field(data_type.value_field))
+    if pa.types.is_map(data_type):
+        return pa.map_(
+            _nullable_type(data_type.key_type),
+            nullable_field(data_type.item_field),
+            keys_sorted=data_type.keys_sorted,
+        )
+    return data_type
+
+
+# The predicate is deliberately conservative: hidden nulls may request a 
check, but a
+# null-free superset proves that all visible values satisfy the required-field 
contract.
+NullCheckPlan = tuple[Callable[[pa.Array], bool], NullChecker]
+
+
+def _null_check_plan(expected_type: pa.DataType) -> Optional[NullCheckPlan]:
+    def field_plan(field: pa.Field) -> Optional[NullCheckPlan]:
+        nested = _null_check_plan(field.type)
+        if field.nullable:
+            return nested
+
+        def needs_check(values: pa.Array) -> bool:
+            return bool(values.null_count) or (nested is not None and 
nested[0](values))
+
+        def check(values: pa.Array) -> None:
+            if values.null_count:
+                raise ValueError(
+                    f"In-process UDF returned nulls in non-nullable field 
{field.name}"
+                )
+            if nested is not None:
+                nested[1](values)
+
+        return needs_check, check
+
+    if pa.types.is_struct(expected_type):
+        fields = [(i, field_plan(f)) for i, f in enumerate(expected_type)]
+        checks = [(i, plan) for i, plan in fields if plan is not None]
+        if not checks:
+            return None
+
+        def needs_struct(array: pa.Array) -> bool:
+            return any(plan[0](array.field(i)) for i, plan in checks)
+
+        def check_struct(array: pa.Array) -> None:
+            valid = None
+            for i, (needs, check) in checks:
+                values = array.field(i)
+                if needs(values):
+                    if array.null_count:
+                        if valid is None:
+                            valid = pc.is_valid(array)
+                        # Filter only the child requiring a check, not its 
sibling payloads.
+                        values = pc.filter(values, valid)
+                    check(values)
+
+        return needs_struct, check_struct
+    if pa.types.is_list(expected_type) or 
pa.types.is_large_list(expected_type):
+        plan = field_plan(expected_type.value_field)
+        if plan is not None:
+
+            def check_list(array: pa.Array) -> None:
+                if plan[0](array.values):
+                    plan[1](pc.list_flatten(array))
+
+            return lambda array: plan[0](array.values), check_list
+    if pa.types.is_map(expected_type):
+        key_plan = _null_check_plan(expected_type.key_type)
+        item_plan = field_plan(expected_type.item_field)
+        # Arrow validation rejects null keys already; only their descendants 
need checks.
+        checks = [(i, p) for i, p in enumerate((key_plan, item_plan)) if p is 
not None]
+        if not checks:
+            return None
+
+        def entries(array: pa.Array) -> pa.Array:
+            start = array.offsets[0].as_py()
+            length = array.offsets[-1].as_py() - start
+            # values.field honors the entries struct's offset; keys/items do 
not.
+            return array.values.slice(start, length)
+
+        def needs_map(array: pa.Array) -> bool:
+            values = entries(array)
+            return any(plan[0](values.field(i)) for i, plan in checks)
+
+        def check_map(array: pa.Array) -> None:
+            if needs_map(array):
+                visible = pc.filter(array, pc.is_valid(array)) if 
array.null_count else array
+                values = entries(visible)
+                for i, (needs, check) in checks:
+                    if needs(values.field(i)):
+                        check(values.field(i))
+
+        return needs_map, check_map
+    return None
+
+
+def _null_checker(expected_type: pa.DataType) -> Optional[NullChecker]:
+    plan = _null_check_plan(expected_type)
+    return plan[1] if plan is not None else None
+
+
+def _has_offset(array: pa.Array) -> bool:
+    if array.offset:
+        return True
+    if pa.types.is_struct(array.type):
+        return any(_has_offset(array.field(i)) for i in 
range(array.type.num_fields))
+    if pa.types.is_list(array.type) or pa.types.is_large_list(array.type):
+        return _has_offset(array.values)
+    if pa.types.is_map(array.type):
+        return _has_offset(array.values)
+    return False
+
+
+def _with_schema(array: pa.Array, expected_type: pa.DataType) -> pa.Array:
+    # Rebind buffers after validating logical nullability. Arrow cast checks 
hidden child
+    # slots too, rejecting null children underneath null parents. from_buffers 
preserves
+    # those masks and applies the declared names, metadata and nullability 
without casting.
+    children = None
+    if pa.types.is_struct(expected_type):
+        children = [_with_schema(array.field(i), f.type) for i, f in 
enumerate(expected_type)]
+    elif pa.types.is_list(expected_type) or 
pa.types.is_large_list(expected_type):
+        children = [_with_schema(array.values, expected_type.value_type)]
+    elif pa.types.is_map(expected_type):
+        entries_type = pa.struct([expected_type.key_field, 
expected_type.item_field])
+        children = [_with_schema(array.values, entries_type)]
+    return pa.Array.from_buffers(
+        expected_type,
+        len(array),
+        array.buffers()[: array.type.num_buffers],
+        null_count=array.null_count,
+        children=children,
+    )
+
+
+def _validate_result(
+    result: pa.Array,
+    expected_rows: int,
+    expected_type: pa.DataType,
+    null_checker: Optional[NullChecker] = None,
+) -> pa.Array:
+    if not isinstance(result, pa.Array):
+        raise TypeError(f"In-process UDF must return a pyarrow.Array, got 
{type(result).__name__}")
+    if len(result) != expected_rows:
+        raise ValueError(f"In-process UDF returned {len(result)} rows; 
expected {expected_rows}")
+    if _nullable_type(result.type) != _nullable_type(expected_type):

Review Comment:
   Timezone-aware results now adopt the expected timezone label without 
changing their UTC instants or copying their buffers. String/binary offset 
widths are converted to the expected representation, including nested values. 
Other type mismatches remain errors, including naive versus aware timestamps 
and differing timestamp units.
   
   Tests cover UTC, Etc/UTC, America/Los_Angeles, and nested string/binary 
results with large types enabled. The guide describes these conversions and 
their allocation behavior.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,344 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.io.File
+import java.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, PyConfig, 
SharedInterpreter}
+import org.apache.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.sql.util.ArrowUtils
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  val SITE_PACKAGES_CONFIG = "spark.inprocess.python.sitePackages"
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var mainConfigured = false
+  @volatile private var sharedConfigured = false
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  private def configureInterpreter(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:
   Kept isolated initialization and added an ASCII-locale warning with the 
`LC_ALL=C.UTF-8` guidance. The guide explicitly states the UTF-8 locale 
requirement and that Python encoding environment flags are ignored.
   
   Standard output and error now use line buffering and write-through, and are 
flushed before interpreter close. A fresh-JVM test with `LC_ALL=C` verifies the 
warning and shutdown flush.



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

Review Comment:
   Documented these as unsupported for this initial implementation: 
`SparkFiles`, `--py-files`, `spark.submit.pyFiles`, `SparkContext.addPyFile`, 
and `SparkSession.addArtifacts(..., pyfile=True)`. The guide explains that 
YARN's process `PYTHONPATH` may make some modules importable incidentally, but 
this is not portable support for those APIs. It directs users to configured 
executor paths for files and preinstalled packages or `sitePackages` for 
modules.



##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala:
##########
@@ -1794,6 +1795,7 @@ object CollapseProject extends Rule[LogicalPlan] with 
AliasHelper {
           lazy val containsUDF = a.child.exists {
             case udf: PythonUDF =>
               isScalarPythonUDF(udf) &&
+                udf.evalType != PythonEvalType.SQL_SCALAR_ARROW_INPROCESS_UDF 
&&

Review Comment:
   Moved the eval-type-258 exclusion to construction of the chaining eval-type 
set, covering both `mergeProjectExpressions` and `canCollapseExpressions`. 
Added the Project-over-Aggregate regression alongside the existing 
Project-over-Project case. All 21 `CollapseProjectSuite` tests pass, including 
existing worker-UDF chaining coverage.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonRuntime.scala:
##########
@@ -0,0 +1,344 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.python
+
+import java.io.File
+import java.nio.ByteBuffer
+import java.util.concurrent.{Callable, ExecutionException, Executors, 
ThreadFactory, TimeoutException, TimeUnit}
+
+import scala.jdk.CollectionConverters._
+
+import jep.{JepConfig, JepException, MainInterpreter, PyConfig, 
SharedInterpreter}
+import org.apache.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.sql.util.ArrowUtils
+import org.apache.spark.util.Utils
+
+/** Owns one interpreter generation per executor plugin lifecycle. */
+private[python] object InProcessPythonRuntime extends Logging {
+  val SITE_PACKAGES_CONFIG = "spark.inprocess.python.sitePackages"
+  private val TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+  private var active: InterpreterSession = _
+  private var mainConfigured = false
+  @volatile private var sharedConfigured = false
+
+  private[python] class LifecycleException(message: String) extends 
IllegalStateException(message)
+
+  private def configureInterpreter(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(new 
JepConfig().addIncludePaths(sitePackages: _*))
+    }
+  }
+
+  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(_bootstrap_error)) from None\n"
+  }
+
+  def initialize(sitePackages: Seq[String] = Seq.empty): Unit = synchronized {
+    if (active != null && !active.isTerminated) {
+      active.requireCompatible(sitePackages)
+    } else {
+      configureInterpreter(sitePackages)
+      val candidate = new InterpreterSession(sitePackages)
+      try {
+        candidate.initialize()
+        active = candidate
+      } catch {
+        case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
candidate.shutdown() }
+      }
+    }
+  }
+
+  def currentSession: InterpreterSession = synchronized {
+    checkState(active != null && active.isRunning)
+    active
+  }
+
+  def shutdown(): Unit = {
+    val session = synchronized { active }
+    if (session != null) session.shutdown()
+  }
+
+  private def checkState(running: Boolean): Unit = {
+    checkState(running, "In-process Python is not running; initialize the 
executor plugin first")
+  }
+
+  private def checkState(running: Boolean, message: String): Unit = {
+    if (!running) throw new IllegalStateException(message)
+  }
+
+  /**
+   * Tasks retain this generation, so stale tasks cannot enter a later 
SparkContext's interpreter.
+   * Lifecycle operations only hold the monitor while enqueueing work, never 
while running Python.
+   */
+  private[python] class InterpreterSession(val sitePackages: Seq[String] = 
Seq.empty) {
+    // CPython native calls need more stack than the usual JVM thread default. 
This is a
+    // platform-dependent size request, not protection against arbitrary 
native crashes.
+    private val executor = Executors.newSingleThreadExecutor(new ThreadFactory 
{
+      override def newThread(runnable: Runnable): Thread = {
+        val thread = new Thread(null, runnable, "inprocess-python", 8L * 1024 
* 1024)
+        thread.setDaemon(true)
+        thread
+      }
+    })
+    @volatile private var running = true
+    // Accessed only on the owning thread.
+    private var interp: SharedInterpreter = _
+
+    def isRunning: Boolean = running
+    def isTerminated: Boolean = executor.isTerminated
+
+    def requireCompatible(paths: Seq[String]): Unit = {
+      if (!isRunning) {
+        throw new LifecycleException("In-process Python is still stopping. 
Wait for outstanding " +
+          "native work to finish or replace the executor process before 
starting a new context.")
+      }
+      if (sitePackages != paths) {
+        throw new LifecycleException("In-process Python is already running 
with different " +
+          "sitePackages. Stop the existing context before changing interpreter 
configuration.")
+      }
+    }
+
+    private[python] def onInterpreterThread[T](body: => T): T = {
+      val context = Option(TaskContext.get())
+      context.foreach(_.killTaskIfInterrupted())
+      val gate = new Object
+      var started = false
+      var cancelled = false
+      val future = synchronized {
+        checkState(running)
+        executor.submit(new Callable[T] {
+          override def call(): T = {
+            gate.synchronized {
+              if (cancelled) throw new TaskKilledException("Cancelled before 
Python invocation")
+              started = true
+            }
+            body
+          }
+        })
+      }
+      var interrupted = false
+      try {
+        while (true) {
+          val taskCancelled = context.exists(_.isInterrupted())
+          if (interrupted || taskCancelled) {
+            val cancelledBeforeStart = gate.synchronized {
+              if (started) false else {
+                cancelled = true
+                future.cancel(false)
+                true
+              }
+            }
+            if (cancelledBeforeStart) {
+              context.foreach(_.killTaskIfInterrupted())
+              throw new InterruptedException("Cancelled before Python 
invocation")
+            }
+          }
+          try {
+            val result = future.get(100, TimeUnit.MILLISECONDS)
+            context.foreach(_.killTaskIfInterrupted())
+            return result
+          } catch {
+            case _: TimeoutException =>
+            case _: InterruptedException => interrupted = true
+            case e: ExecutionException => throw e.getCause
+          }
+        }
+        throw new IllegalStateException("Unreachable")
+      } finally {
+        // Once native work starts, wait for it even after cancellation: the 
caller still owns
+        // CDI structs that Python may use. Pending work, however, is safe to 
cancel immediately.
+        if (interrupted) Thread.currentThread().interrupt()
+      }
+    }
+
+    def initialize(): Unit = onInterpreterThread {
+      val candidate = new ManagedSharedInterpreter()
+      try {
+        candidate.set("_site_packages", sitePackages.asJava)
+        val sparkPaths = PythonUtils.mergePythonPaths(
+          PythonUtils.sparkPythonPath, sys.env.getOrElse("PYTHONPATH", ""))
+          .split(File.pathSeparator).filter(_.nonEmpty)
+        candidate.set("_spark_paths", sparkPaths.toSeq.asJava)
+        candidate.exec(bootstrapScript(
+          """import os, site, sys
+            |_configured = [os.path.abspath(p) for p in _site_packages]
+            |_before = set(sys.path)
+            |for _path in _configured:
+            |    site.addsitedir(_path)
+            |_added = [p for p in sys.path if p not in _before and p not in 
_configured]
+            |_preferred = list(dict.fromkeys(list(_spark_paths) + _configured 
+ _added))
+            |sys.path[:] = _preferred + [p for p in sys.path if p not in 
_preferred]
+            |del _site_packages, _spark_paths, _configured, _before, _added, 
_preferred
+            |""".stripMargin))
+        candidate.exec(bootstrapScript(
+          "from pyspark.sql.pandas.utils import 
require_minimum_pyarrow_version\n" +
+          "require_minimum_pyarrow_version()\n" +
+          "from pyspark.inprocess.runtime import " +
+          "_inprocess_invoke, _inprocess_register, _inprocess_release, _udfs"))
+        interp = candidate
+      } catch {
+        case t: Throwable => Utils.tryWithSafeFinally { throw t } { 
candidate.close() }
+      }
+    }
+
+    /** Enqueue cleanup after outstanding calls without creating an executor 
or waiting. */
+    def release(handles: Seq[String]): Unit = synchronized {
+      if (running && handles.nonEmpty) {
+        executor.submit(new Runnable {
+          override def run(): Unit = {
+            if (interp != null) interp.invoke("_inprocess_release", 
handles.asJava)
+          }
+        })
+      }
+      // During shutdown the queued close clears all remaining handles.
+    }
+
+    /** A timeout bounds plugin stop, not native execution or CDI buffer 
ownership. */
+    def shutdown(waitMillis: Long = 5000L): Unit = {
+      synchronized {
+        if (running) {
+          running = false
+          executor.submit(new Runnable {
+            override def run(): Unit = {
+              if (interp != null) {
+                try {
+                  interp.exec("_udfs.clear()")
+                } finally {
+                  try { interp.close() } finally { interp = null }
+                }
+              }
+            }
+          })
+          executor.shutdown()
+        }
+      }
+      try {
+        if (!executor.awaitTermination(waitMillis, TimeUnit.MILLISECONDS)) {
+          logWarning("In-process Python is still stopping; native work and its 
buffers " +
+            "remain alive until the invocation finishes or the process exits.")
+        }
+      } catch {
+        case _: InterruptedException => Thread.currentThread().interrupt()
+      }
+    }
+
+    private[python] def timedOnInterpreterThread(body: => Unit): Long = 
onInterpreterThread {
+      val start = System.nanoTime()
+      body
+      System.nanoTime() - start
+    }
+
+    def register(
+        handle: String,
+        serializedUdf: Array[Byte],
+        expectedField: Field,
+        pythonVersion: String,
+        hideTraceback: Boolean,
+        simplifiedTraceback: Boolean,
+        tracebackWithLocals: Boolean): Long = {
+      // Bulk-copy on the task thread. JEP's PyJBuffer supports memoryview 
without per-byte JNI.
+      val command = ByteBuffer.allocateDirect(serializedUdf.length)

Review Comment:
   The temporary command copy now uses `ArrowUtils.rootAllocator`, passing its 
`nioBuffer` to JEP and closing the allocation in `finally`. It no longer 
depends on GC for native-memory reclamation. Added an allocation-balance 
assertion for registration failure; successful registration and large closures 
remain covered by the integration tests.



##########
python/pyspark/inprocess/runtime.py:
##########
@@ -0,0 +1,310 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+
+"""Arrow CDI entry points called on the executor's dedicated JEP interpreter 
thread.
+
+Functions are registered once per task and released when that task finishes. 
Calls
+pass only a handle and CDI addresses, so large closures are not copied per 
batch.
+"""
+
+import sys
+from typing import Any, Callable, Iterable, Optional, Sequence
+
+import pyarrow as pa
+import pyarrow.compute as pc
+
+from pyspark import cloudpickle
+from pyspark.errors import PySparkRuntimeError
+from pyspark.sql.pandas.utils import require_minimum_pyarrow_version
+from pyspark.util import _format_exception
+
+_UDF_TRACEBACK_SENTINEL = "__INPROCESS_UDF_TRACEBACK__:"
+NullChecker = Callable[[pa.Array], None]
+_udfs: dict[str, tuple[Callable[..., pa.Array], pa.DataType, NullChecker, 
bool, bool, bool]] = {}
+
+
+def _jep_safe_message(message: str) -> str:
+    # JEP uses JNI modified UTF-8 for exception text. Keep the transport ASCII 
and
+    # escape NUL explicitly; ordinary UTF-8 and embedded NUL are not safe here.
+    return message.encode("ascii", 
"backslashreplace").decode("ascii").replace("\0", "\\x00")
+
+
+def _inprocess_register(
+    handle: str,
+    serialized_udf: Any,
+    schema_ptr: int,
+    python_version: str,
+    hide_traceback: bool = False,
+    simplified_traceback: bool = False,
+    traceback_with_locals: bool = False,
+) -> None:
+    try:
+        require_minimum_pyarrow_version()
+        embedded_version = "%d.%d" % sys.version_info[:2]
+        if python_version != embedded_version:
+            raise PySparkRuntimeError(
+                errorClass="PYTHON_VERSION_MISMATCH",
+                messageParameters={
+                    "worker_version": embedded_version,
+                    "driver_version": python_version,
+                },
+            )
+        # JEP exposes direct ByteBuffers through the buffer protocol. Unpickle 
a separate
+        # function per task without iterating over a PyJArray one JNI call per 
byte.
+        func = cloudpickle.loads(memoryview(serialized_udf))
+        # The JVM is the single source of truth for Arrow layout and logical 
metadata.
+        expected_type = pa.Field._import_from_c(schema_ptr).type
+        checker = _null_checker(expected_type) or (lambda array: None)
+        _udfs[handle] = (
+            func,
+            expected_type,
+            checker,
+            hide_traceback,
+            simplified_traceback,
+            traceback_with_locals,
+        )
+    except BaseException as error:
+        # In JEP, an uncaught SystemExit can terminate the entire executor JVM.
+        raise RuntimeError(
+            _UDF_TRACEBACK_SENTINEL
+            + _jep_safe_message(
+                _format_exception(
+                    error, hide_traceback, simplified_traceback, 
traceback_with_locals
+                )
+            )
+        ) from None
+
+
+def _inprocess_release(handles: Iterable[str]) -> None:
+    for handle in handles:
+        _udfs.pop(handle, None)
+
+
+def _nullable_type(data_type: pa.DataType) -> pa.DataType:
+    def nullable_field(field: pa.Field) -> pa.Field:
+        return pa.field(field.name, _nullable_type(field.type), nullable=True)
+
+    if pa.types.is_struct(data_type):
+        return pa.struct([nullable_field(field) for field in data_type])
+    if pa.types.is_list(data_type):
+        return pa.list_(nullable_field(data_type.value_field))
+    if pa.types.is_large_list(data_type):
+        return pa.large_list(nullable_field(data_type.value_field))
+    if pa.types.is_map(data_type):
+        return pa.map_(
+            _nullable_type(data_type.key_type),
+            nullable_field(data_type.item_field),
+            keys_sorted=data_type.keys_sorted,
+        )
+    return data_type
+
+
+# The predicate is deliberately conservative: hidden nulls may request a 
check, but a
+# null-free superset proves that all visible values satisfy the required-field 
contract.
+NullCheckPlan = tuple[Callable[[pa.Array], bool], NullChecker]
+
+
+def _null_check_plan(expected_type: pa.DataType) -> Optional[NullCheckPlan]:
+    def field_plan(field: pa.Field) -> Optional[NullCheckPlan]:
+        nested = _null_check_plan(field.type)
+        if field.nullable:
+            return nested
+
+        def needs_check(values: pa.Array) -> bool:
+            return bool(values.null_count) or (nested is not None and 
nested[0](values))
+
+        def check(values: pa.Array) -> None:
+            if values.null_count:
+                raise ValueError(
+                    f"In-process UDF returned nulls in non-nullable field 
{field.name}"
+                )
+            if nested is not None:
+                nested[1](values)
+
+        return needs_check, check
+
+    if pa.types.is_struct(expected_type):
+        fields = [(i, field_plan(f)) for i, f in enumerate(expected_type)]
+        checks = [(i, plan) for i, plan in fields if plan is not None]
+        if not checks:
+            return None
+
+        def needs_struct(array: pa.Array) -> bool:
+            return any(plan[0](array.field(i)) for i, plan in checks)
+
+        def check_struct(array: pa.Array) -> None:
+            valid = None
+            for i, (needs, check) in checks:
+                values = array.field(i)
+                if needs(values):
+                    if array.null_count:
+                        if valid is None:
+                            valid = pc.is_valid(array)
+                        # Filter only the child requiring a check, not its 
sibling payloads.
+                        values = pc.filter(values, valid)
+                    check(values)
+
+        return needs_struct, check_struct
+    if pa.types.is_list(expected_type) or 
pa.types.is_large_list(expected_type):
+        plan = field_plan(expected_type.value_field)
+        if plan is not None:
+
+            def check_list(array: pa.Array) -> None:
+                if plan[0](array.values):
+                    plan[1](pc.list_flatten(array))
+
+            return lambda array: plan[0](array.values), check_list
+    if pa.types.is_map(expected_type):
+        key_plan = _null_check_plan(expected_type.key_type)
+        item_plan = field_plan(expected_type.item_field)
+        # Arrow validation rejects null keys already; only their descendants 
need checks.
+        checks = [(i, p) for i, p in enumerate((key_plan, item_plan)) if p is 
not None]
+        if not checks:
+            return None
+
+        def entries(array: pa.Array) -> pa.Array:
+            start = array.offsets[0].as_py()

Review Comment:
   Added the zero-length guard before reading map offsets. The regression test 
uses an empty-map stub whose offsets accessor fails if touched, plus normal 
empty nested maps, to verify that the offsets read is skipped. I did not 
reproduce the native crash.



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