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]
