dongjoon-hyun commented on code in PR #58978: URL: https://github.com/apache/spark/pull/58978#discussion_r4117400983
########## 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: **[Medium] Result vectors backed by Python-owned buffers are released on Spark task threads, so Python runs there and can block on the GIL.** `closeBatch()` (and `close()` from `hasNext` or the task completion listener) closes the vectors imported by `InProcessArrowBridge.cdiToColumn` on the task thread. When the last buffer is released, Arrow Java 19 calls the exported array's C `release` callback synchronously on that thread (`ReferenceCountedArrowArray.release()` -> `array.release()`). If the result buffers are owned by Python objects, PyArrow 18's destructors then take the GIL on the task thread: `NumPyBuffer::~NumPyBuffer() { PyAcquireGIL lock; Py_XDECREF(arr_); }`, and `PyBuffer::~PyBuffer()` likewise. `_with_schema` keeps the original buffers, so this is the normal case for zero-copy NumPy results, e.g. the `random_noise` example in the `inprocess_udf` docstring (`pa.array(np.random.randint(0, 100, len(x)), type=pa.int64())`), `pa.array(np.log1p(x.to_numpy()))` or `pa.Array.from_pandas(series)`. When an executor runs more than one task (supported, and exercised by the tests with `local[2]` and `spark.task.cpus=0.5`): - Task A's cleanup waits in native `PyGILState_Ensure` while task B's UDF holds the GIL in a long native call (a backtracking `re` match, a C extension that keeps the GIL) or hangs. The wait cannot be interrupted, so A can neither finish nor be killed, although its Python work is done. - NumPy deallocation, `__del__` and weakref callbacks run on arbitrary task threads with the default JVM stack, outside the single interpreter thread model described in the guide. This is not a deadlock (the waiting thread only holds the monitor of its own foreign allocation), but it breaks the threading contract. One option: have `_inprocess_invoke` keep a reference to each exported result and drop it on the next `_inprocess_invoke` or `_inprocess_release` for that handle. Then the JVM's release only decrements a shared pointer, and the last Python reference is always dropped on the interpreter thread. ########## 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: **[Low, performance] Follow-up: this exemption covers `mergeProjectExpressions`, but not `canCollapseExpressions`.** Thanks for adding this. However, the `Project` over `Aggregate` branch of `traverse` (L1702-1711) goes through `canCollapseExpressions`, whose `inline` (L1858-1863) still treats 258 like a chaining eval type, and L1677-1679 still add 258 to `pythonUDFEvalTypesInUpperProjects`. For example: ```python agg = df.groupBy("k").agg(ip(F.sum("v")).alias("s")) agg.withColumn("t", ip2("s")) # or agg.select(ip2("s"), "s") ``` `inline = true` skips the reference-count and cheapness checks, and the plan becomes `Aggregate [k, ip(sum(v)) AS s, ip2(ip(sum(v))) AS t]`. `ExtractPythonUDFFromAggregate` then gives each `sum(v)` occurrence its own alias, so `ip(agg#1)` and `ip(agg#2)` are not deduplicated and `ip` runs twice per group, with the same number of in-process nodes as without the collapse. The new `CollapseProjectSuite` test only covers `Project` over `Project`. Suggestion: exclude 258 where the eval type set is built (L1677-1679), which covers both consumers, and add a test like `testRelation.groupBy($"a")(udf(sum($"b")).as("s")).select(udf($"s").as("r"), $"s")` that stays unchanged after optimization. ########## 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: **[Medium] JEP's Java import hook is installed with the default `ClassList`, so Python packages named like a Java package on the executor classpath cannot be imported.** `new JepConfig()` has no `ClassEnquirer`. When the first `SharedInterpreter` is configured, `Jep.configureInterpreter` calls `setupJavaImportHook(null)`, which falls back to `ClassList.getInstance()`, and `java_import_hook.py` inserts `JepJavaImporter` at `sys.meta_path[0]`. Its `find_spec` claims every name for which `isJavaPackage(fullname)` is true and returns an empty Java package module (`__path__ = []`, `__file__ = '<java>'`) without falling back to Python. Only `io`, `re` and `py4j` are excluded. Examples that work in worker UDFs but fail here: - `from spire.xls import *` (Spire.XLS for Python) fails with `ModuleNotFoundError: No module named 'spire.xls'`, because Spark's own jars contain the Scala `spire` package (an MLlib/breeze dependency). - With Jedis on the system classpath (e.g. spark-redis in the image), `import redis` resolves to the Java `redis.clients` package instead of redis-py. Since the runtime never imports Java from Python, the hook only adds cost: `ClassList` eagerly indexes every class on `java.class.path` (directories are walked recursively; on YARN `{{PWD}}` is on the classpath, so the container directory including an unpacked `--archives` venv is scanned) and keeps the index for the JVM lifetime, and every later Python import makes a JNI `isJavaPackage` call. Two smaller issues on this line: - With the default empty `sitePackages`, `addIncludePaths()` still creates an empty include path, so JEP runs `sys.path += ''.split(':')` and the executor's working directory ends up on `sys.path`, despite the isolated config. - The paths are spliced unescaped into Python source (`"sys.path += '" + includePath + "'.split(...)"`), so a path containing `'` fails plugin init with a `SyntaxError`, and a path containing `:` is split in two. Suggestion: `new JepConfig().setClassEnquirer(new NamingConventionClassEnquirer(false))` (no Java package names, no classpath scan), and call `addIncludePaths` only when `sitePackages.nonEmpty`, rejecting entries that contain `'`, a newline or `File.pathSeparator`. ########## 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: **[Medium] The result type must match the session time zone and `useLargeVarTypes` exactly, which the UDF cannot observe.** The expected type now comes from the JVM (`ArrowUtils.toArrowField("result", udf.dataType, true, timeZoneId, largeVarTypes)` with `conf.sessionLocalTimeZone` and `conf.arrowUseLargeVarTypes`), and this comparison only relaxes nested nullability. The embedded interpreter has no access to the session configuration, so the declared Spark type alone doesn't tell the UDF which Arrow type to produce. Examples, reproduced with `_validate_result`: - `@inprocess_udf("timestamp") def f(s): return pc.assume_timezone(pc.strptime(s, format="%Y-%m-%d", unit="us"), "UTC")` passes only when the session time zone is exactly `UTC`. With the common JVM default `Etc/UTC`, or `America/Los_Angeles`, every task fails with `TypeError: In-process UDF returned timestamp[us, tz=UTC]; expected timestamp[us, tz=Etc/UTC]`. - With `spark.sql.execution.arrow.useLargeVarTypes=true`, a UDF returning `pc.cast(x, pa.string())` fails with `returned string; expected large_string`; binary and nested strings behave the same. The same functions work as `arrow_udf`, because the worker casts the result in `enforce_schema` (`arrow_cast=True`). For `TimestampType` the tz label is only metadata (the values are UTC microseconds), so a tz-aware result could be relabeled with the expected type in `_with_schema` without a copy (no JVM change is needed, since the exported type then equals the expected one), and string/binary could be cast to the large variants. At minimum, please document that the expected tz is `spark.sql.session.timeZone` and that `useLargeVarTypes` changes the expected string/binary types. ########## 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: **[Medium] `PyConfig.isolated()` also turns off UTF-8 mode, locale coercion and unbuffered stdio, which Python workers rely on.** Switching to the isolated config addressed the fatal signal handlers, but JEP calls `PyConfig_InitIsolatedConfig` and then `Py_InitializeFromConfig`, so CPython derives an isolated pre-config as well (`PyPreConfig_InitIsolatedConfig`: `utf8_mode=0`, `configure_locale=0`, `coerce_c_locale=0`, `use_environment=0`). Two visible differences from worker UDFs: 1. On executors without a UTF-8 locale (no `LANG`/`LC_ALL`, e.g. `ubuntu`-based images such as the python-312 CI image, or YARN NodeManagers started without `LANG`), the embedded interpreter uses ASCII: `print("café")` raises `UnicodeEncodeError`, and `open(path).read()` of a UTF-8 file raises `UnicodeDecodeError`. Worker processes in the same environment get PEP 538/540 behavior and work. `PYTHONUTF8` and `PYTHONIOENCODING` set through `spark.executorEnv.*` are ignored; only `LC_ALL=C.UTF-8` (or `LANG`) helps. 2. `buffered_stdio` keeps its default and `PYTHONUNBUFFERED` is ignored, so `sys.stdout` on the executor's fd 1 (a file or pipe) is block-buffered. JEP never calls `Py_Finalize` (only from `JNI_OnUnload`), and `InterpreterSession.shutdown` doesn't flush, so UDF `print()` output appears late and the last buffer is lost when the executor exits. Workers run with `PYTHONUNBUFFERED=YES`. Suggestion: `sys.stdout.reconfigure(write_through=True)` (or line buffering) in the bootstrap and a flush of `sys.stdout`/`sys.stderr` before `interp.close()`; document the UTF-8 locale requirement with the `LC_ALL=C.UTF-8` workaround, or warn at bootstrap when the locale encoding is ASCII. ########## 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: **[Low] `entries()` can segfault the executor for a zero-length map without an offsets buffer.** Arrow allows a zero-length list/map array whose offsets buffer is null (ARROW-544), and both `validate()` and `validate(full=True)` accept it. `array.offsets` then wraps the null buffer as a length-1 array, and `offsets[0]` reads through a null pointer (Arrow 18's `BoxOffsets` doesn't guard it either). Since Python runs inside the executor JVM, this is a SIGSEGV of the executor, and a deterministic one, so every retry loses another executor. Example: return type `ArrayType(MapType(StringType(), IntegerType(), valueContainsNull=False))`, and a result whose lists are all empty and whose map child was built with `pa.Array.from_buffers(map_t, 0, [None, None], children=[entries])`. `_validate_result` -> `check_list` -> `needs_map` -> `entries` exits with code 139 (reproduced with pyarrow 15). Standard PyArrow APIs and the CDI import path don't produce such arrays, so this takes hand-built buffers, but without the crash it would have been a clean task error (Arrow Java's importer rejects the null buffer). Suggestion: `if len(array) == 0: return array.values.slice(0, 0)` before reading the offsets. ########## sql/core/src/main/scala/org/apache/spark/sql/execution/python/InProcessPythonUDFBuilder.scala: ########## @@ -0,0 +1,92 @@ +/* + * 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.{Collections, List => JList} + +import scala.jdk.CollectionConverters._ + +import org.apache.spark.{SparkEnv, SparkException} +import org.apache.spark.api.python.{PythonEvalType, SimplePythonFunction} +import org.apache.spark.sql.Column +import org.apache.spark.sql.catalyst.expressions.PythonUDF +import org.apache.spark.sql.catalyst.plans.logical.NamedParametersSupport +import org.apache.spark.sql.classic.{ColumnNodeExpression, ExpressionUtils} +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types.DataType + +/** + * JVM-side builder for in-process [[PythonUDF]] expressions, called from the Python API + * via py4j's JVM reflection bridge (``sc._jvm.org.apache.spark...InProcessPythonUDFBuilder``). + * + * Accepts Java-typed arguments as passed by PySpark's ``sc._jvm`` proxy and returns a + * [[Column]] backed by a [[PythonUDF]] with the in-process evaluation type. + */ +object InProcessPythonUDFBuilder { + + /** + * Build a [[Column]] backed by an in-process [[PythonUDF]] expression. + * + * @param name display name (Python function ``__name__``) + * @param serializedFunc cloudpickle bytes of the Python UDF + * @param returnTypeJson JSON string of the Spark SQL return type + * @param jColumns Java List of JVM [[Column]] objects (the UDF inputs) + * @param deterministic whether the UDF always returns the same output for the same input; + * set to false for UDFs that use randomness or external state + * @param pythonVersion driver's Python major.minor version + * @return [[Column]] backed by an in-process [[PythonUDF]] expression + */ + def build( + name: String, + serializedFunc: Array[Byte], + returnTypeJson: String, + jColumns: JList[Column], + deterministic: Boolean, + pythonVersion: String): Column = { + checkConfiguration(SQLConf.get) Review Comment: **[Low] The build-time check reads the py4j thread's active session, not the session the Column runs in.** `SQLConf.get` here resolves to the JVM active session of the py4j thread (or an empty fallback conf if there is none). A `Column` is session-agnostic, and the authoritative check with the plan's session already runs in `InProcessArrowEvalPythonExec.doExecute`. - False rejection: after `spark.conf.set("spark.sql.pyspark.udf.profiler", "perf")` (or a `spark.pythonWorkerEnv.*` entry), `spark2 = spark.newSession(); df2 = spark2.range(3); df2.select(f(df2.id))` fails with `INVALID_SPARK_CONFIG.UNSUPPORTED_IN_PROCESS_PYTHON_UDF`, although `spark2` has neither setting, because PySpark's `newSession()` doesn't change the JVM active session. - No-op: on non-main Python threads (pinned thread mode), the py4j thread has no active session, so `SQLConf.get` is a fresh fallback conf and nothing is checked. Suggestion: keep only the application-level `spark.executor.pyspark.memory` check here and rely on `doExecute` for session configurations. Separately, other worker-only settings such as `spark.sql.pyspark.worker.logging.enabled` and `spark.sql.execution.pyspark.udf.faulthandler.enabled` are neither rejected nor documented, so enabling them silently has no effect on in-process UDFs. ########## 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: **[Low] The direct buffer that holds the pickled closure is never freed explicitly.** Every registration (once per task per UDF) allocates a direct `ByteBuffer` as large as the whole command, and only GC releases it. JEP drops its reference when the call returns (`Py_CLEAR(pyargs)` -> `DeleteGlobalRef`) and `cloudpickle.loads` copies the data, so the buffer is dead right after `_inprocess_register`. Its native memory, however, is freed only when a GC collects the small `DirectByteBuffer` object. If that object was tenured while the registration waited behind other tasks on the interpreter thread, or during a long unpickle, this can take until the next old-gen or concurrent cycle. In-process UDFs reject broadcasts and have no broadcast threshold, so a captured model is shipped inline at any size. With a closure of a few hundred MB and several tasks per executor, several untracked off-heap copies can stay alive (a jshell reproduction kept 4 x 128 MB across 100 young GCs), which can push RSS beyond `memoryOverhead`, or fail with `OutOfMemoryError: Cannot reserve ... direct buffer memory` when explicit GC is disabled. Suggestion: free it in the `finally`, e.g. with `StorageUtils.dispose(command)`, or allocate it from `ArrowUtils.rootAllocator` and pass `nioBuffer(0, n)` so that it can be closed. The extra `func.command.toArray` heap copy in the caller could be avoided as well. ########## docs/sql-pyspark-inprocess-udf.md: ########## @@ -0,0 +1,658 @@ +--- +layout: global +title: In-Process Python UDFs Review Comment: **[Low, docs] This page isn't linked from anywhere.** Neither `docs/_data/menu-sql.yaml` (which lists `sql-pyspark-pandas-with-arrow.html`) nor any other page links to `sql-pyspark-inprocess-udf.html`, and `pyspark.inprocess.inprocess_udf` is not in the PySpark API reference under `python/docs/source`. The new sentence in `docs/configuration.md` doesn't link here either. Without knowing the URL, users can't find the setup requirements and the caveats (a hung call blocks the executor, a native crash kills it). Suggestion: add a menu entry next to the Arrow guide, link it from `configuration.md` and the UDF pages, and add the API to the reference docs. ########## python/pyspark/inprocess/udf.py: ########## @@ -0,0 +1,232 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +""" +Python API for in-process UDF registration. + +Usage:: + + import pyarrow.compute as pc + from pyspark.inprocess import inprocess_udf + from pyspark.sql.types import LongType + + @inprocess_udf(return_type=LongType()) + def double(x): + # x is a pa.Array; return a pa.Array + return pc.multiply(x, 2) + + df.select(double(df.value)).show() +""" + +import io +import sys +from functools import update_wrapper +from inspect import signature +from typing import Any, Callable, Optional, Union + +from pyspark import Accumulator, Broadcast, cloudpickle +from pyspark.errors import PySparkTypeError, PySparkValueError +from pyspark.sql.column import Column +from pyspark.sql.types import DataType, _parse_datatype_string +from pyspark.util import PythonEvalType + + +class _InProcessPickler(cloudpickle.CloudPickler): + def reducer_override(self, obj: Any) -> Any: + if isinstance(obj, (Broadcast, Accumulator)): + raise TypeError("In-process UDFs do not support Spark broadcasts or accumulators") + return super().reducer_override(obj) + + +def _serialize_udf(func: Callable) -> bytes: + buffer = io.BytesIO() + _InProcessPickler(buffer).dump(func) + return buffer.getvalue() + + +class InProcessUDFWrapper: + """ + Wraps a Python function as an in-process UDF. + + Returned by ``@inprocess_udf``. Calling an instance with Spark ``Column`` + arguments creates a ``Column`` expression backed by ``PythonUDF`` + on the JVM side. + """ + + def __init__( + self, func: Callable, return_type: Union[DataType, str], deterministic: bool = True + ) -> None: + if not isinstance(return_type, (DataType, str)): + raise PySparkTypeError( + errorClass="NOT_EXPECTED_TYPE", + messageParameters={ + "expected_type": "DataType or str", + "arg_name": "return_type", + "arg_type": type(return_type).__name__, + }, + ) + self._return_type = return_type + self._parsed_return_type: Optional[DataType] = None + self.evalType = PythonEvalType.SQL_SCALAR_ARROW_INPROCESS_UDF + self._deterministic: bool = deterministic + self._name: str = getattr(func, "__name__", "inprocess_udf") + + from pyspark.sql.pandas.utils import require_minimum_pyarrow_version + + require_minimum_pyarrow_version() + if not signature(func).parameters: + raise PySparkValueError( + errorClass="INVALID_PANDAS_UDF", + messageParameters={"detail": "0-arg inprocess_udfs are not supported."}, + ) + self._func = func + self._serialized: Optional[bytes] = None + update_wrapper(self, func, updated=()) + + @property + def func(self) -> Callable: + return self._func + + @property + def returnType(self) -> DataType: + if self._parsed_return_type is None: + parsed = ( + _parse_datatype_string(self._return_type) + if isinstance(self._return_type, str) + else self._return_type + ) + from pyspark.sql.udf import UserDefinedFunction + + UserDefinedFunction._check_return_type(parsed, PythonEvalType.SQL_SCALAR_ARROW_UDF) + from pyspark.sql.pandas.types import to_arrow_type + + to_arrow_type(parsed, timezone="UTC", error_on_duplicated_field_names_in_struct=True) + self._parsed_return_type = parsed + return self._parsed_return_type + + @property + def deterministic(self) -> bool: + return self._deterministic + + def asNondeterministic(self) -> "InProcessUDFWrapper": + self._deterministic = False + return self + + def _serialize(self) -> bytes: + if self._serialized is None: + # Validate before caching the command, including driver-only UDT definitions. + self.returnType + self._serialized = _serialize_udf(self._func) + return self._serialized + + def __call__(self, *cols: Union[Column, str], **kwargs: Union[Column, str]) -> Column: + """ + Create a ``Column`` expression invoking this UDF with the given columns. + + Args: + *cols: Spark ``Column`` objects (e.g. ``df.value``, ``col("x")``) + + Returns: + pyspark.sql.Column + """ + from pyspark import SparkContext + from pyspark.sql.classic.column import _to_java_column + + sc = SparkContext._active_spark_context Review Comment: **[Low] Calling an in-process UDF under Spark Connect fails with misleading errors.** `__call__` has no `is_remote()` guard and reads `SparkContext._active_spark_context` directly: - Remote Connect (`sc://...`): `RuntimeError: No active SparkContext. Start a SparkSession before calling an inprocess_udf.`, although a session exists. - Local Connect (`.remote("local[*]")` or `spark.api.mode=connect`, which creates a classic `SparkContext` in the same process): `f(df.id)` fails with `NOT_EXPECTED_TYPE` ("should be Column or str, got Column"), and `df.select(f("id"))` builds a classic JVM Column that the Connect DataFrame later rejects with `TypeError: 'Column' object is not callable`. The guide says that client SQL registration and server planning reject Connect, and the tests cover only those two paths. Also, `python/pyspark/inprocess` is outside `TARGET_PATHS` in `dev/check_pyspark_custom_errors.py`, so the bare `RuntimeError` / `TypeError` in this module are not caught by that lint. Suggestion: raise `PySparkNotImplementedError` when `is_remote()`, use `get_active_spark_context()` (`SESSION_OR_CONTEXT_NOT_EXISTS`) otherwise, and add a Connect test for the DataFrame API path. ########## 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) + command.put(serializedUdf).flip() + val schema = ArrowSchema.allocateNew(ArrowUtils.rootAllocator) + Utils.tryWithSafeFinally { + Data.exportField(ArrowUtils.rootAllocator, expectedField, null, schema) + timedOnInterpreterThread { + withPythonException { + interp.invoke("_inprocess_register", handle, command, + java.lang.Long.valueOf(schema.memoryAddress()), pythonVersion, + java.lang.Boolean.valueOf(hideTraceback), + java.lang.Boolean.valueOf(simplifiedTraceback), + java.lang.Boolean.valueOf(tracebackWithLocals)) + } + } + } { + Utils.tryWithSafeFinally { + if (schema.snapshot().release != 0L) schema.release() + } { schema.close() } + } + } + + def invoke( + handle: String, + inputArrayPtrs: Array[Long], + inputSchemaPtrs: Array[Long], + outputArrayAddr: Long, + outputSchemaAddr: Long, + expectedRows: Int, + argumentNames: Array[String]): Long = timedOnInterpreterThread { + val arrayPtrs = inputArrayPtrs.map(java.lang.Long.valueOf).toSeq.asJava + val schemaPtrs = inputSchemaPtrs.map(java.lang.Long.valueOf).toSeq.asJava + withPythonException { + interp.invoke("_inprocess_invoke", handle, arrayPtrs, schemaPtrs, + java.lang.Long.valueOf(outputArrayAddr), java.lang.Long.valueOf(outputSchemaAddr), + java.lang.Integer.valueOf(expectedRows), argumentNames.toSeq.asJava) + } + } + } + + private def withPythonException(body: => Unit): Unit = { + try { + body + } catch { + case e: JepException => Review Comment: **[Low] Without the JEP JAR, tasks fail with `NoClassDefFoundError: jep/JepException` instead of the intended plugin error.** `withPythonException` is in `object InProcessPythonRuntime` and catches `JepException` by type. Scala emits a typed exception-table entry for this, and HotSpot's verifier resolves catch types when it links the class (`verify_exception_handler_table` -> `is_assignable_from`), so `InProcessPythonRuntime$` can't be linked at all without the JEP JAR. Scenario: a stock deployment (JEP and `arrow-c-data` are not bundled), no `spark.plugins`, then `spark.range(1).select(inprocess_udf("long")(lambda x: x)("id")).collect()`. Every task attempt fails at `InProcessArrowEvalPythonEvaluatorFactory.scala:137` (`InProcessPythonRuntime.currentSession`) with `java.lang.NoClassDefFoundError: jep/JepException` and is retried until `spark.task.maxFailures`. The "initialize the executor plugin first" message is only reachable when JEP happens to be on the classpath, and nothing on the driver checks that the plugin is registered, although the guide says that "Without this plugin, in-process UDF execution fails with an initialization error". Suggestion: move the typed catch out of the object (e.g. into `InterpreterSession`) or use a guard such as `case e: RuntimeException if e.isInstanceOf[JepException]`, translate a `LinkageError` around `currentSession` into the plugin/JAR hint, and optionally fail fast on the driver when `spark.plugins` doesn't contain `InProcessPythonPlugin`. ########## 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") Review Comment: **[Low] This escapes more than JNI needs, so non-English error messages become unreadable.** Following up on my earlier comment: JEP builds the message with `PyUnicode_AsUTF8` + `NewStringUTF`, and modified UTF-8 is identical to standard UTF-8 for every BMP code point except NUL, while lone surrogates can't be encoded at all. So only NUL, surrogates and supplementary characters (U+10000 and above) actually need escaping. Escaping everything turns `raise ValueError("잘못된 값: café")` into `ValueError: 잘못된 값: caf\xe9` on the driver, escapes non-ASCII source lines and file paths in tracebacks (which also shifts the `^^^^` markers), and since backslashes themselves aren't escaped, the original text can't be recovered. Worker UDFs show the text as is. Suggestion: escape only the unsafe subset, e.g. `re.sub("[\x00\ud800-\udfff\U00010000-\U0010ffff]", lambda m: m.group().encode("unicode_escape").decode("ascii"), message)`. `test_exception_text_is_safe_for_jni` and `test_exception_unicode_and_nul_survive_jep` still pass with that; the guide sentence and this comment would need a small update. ########## docs/sql-pyspark-inprocess-udf.md: ########## @@ -0,0 +1,658 @@ +--- +layout: global +title: In-Process Python UDFs +displayTitle: In-Process Python UDFs +license: | + Licensed to the Apache Software Foundation (ASF) under one or more + contributor license agreements. See the NOTICE file distributed with + this work for additional information regarding copyright ownership. + The ASF licenses this file to You under the Apache License, Version 2.0 + (the "License"); you may not use this file except in compliance with + the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +--- + +* Table of contents +{:toc} + +## Runtime and result contract + +Each executor owns a dedicated interpreter thread. The plugin initializes the +interpreter on that thread, and task calls and shutdown are dispatched to the +same thread. The JVM is asked to allocate an 8 MiB stack for this thread; the +actual size is platform-dependent. Calls from concurrent tasks are queued on the +interpreter thread. +One task per executor is recommended for throughput, but is not a correctness requirement. +Application-level Python parallelism comes from multiple executor JVMs. +The plugin configures JEP's process-wide interpreter with hash seed `0`, matching +Spark's default Python worker seed. It must initialize before any other JEP user in +the JVM. The seed cannot change between SparkContexts in the same process; a custom +worker `PYTHONHASHSEED` does not override this embedded-runtime setting. + +Task cancellation cannot safely stop arbitrary native Python code. An interrupted +caller waits for the current invocation to finish before freeing the Arrow CDI +structures, then restores its interrupt status. A UDF that never returns can +therefore prevent its task from completing cancellation and block every subsequent +in-process UDF on that executor, including calls from other tasks, jobs, and sessions. +Recovery from a permanently hung invocation requires replacing the executor process. +Plugin shutdown stops accepting new calls and waits up to five seconds for the interpreter thread. If a call is +still running, cleanup stays queued behind it; its memory remains live until the +call returns or the process exits. Shutdown does not forcibly interrupt native +code. A new interpreter cannot start until the previous one has fully stopped. + +A scalar UDF must return a `pyarrow.Array` with exactly one element per input row. +The runtime checks the result type against the declared Spark type, including +nested fields, decimal scale, and timestamp unit/timezone. Value types must match +exactly: use an explicit PyArrow cast in the UDF for numeric or other conversions. +Nested field nullability may differ if the actual values satisfy the declared nullability. Sliced results, including nested +child slices, are copied to remove offsets that Arrow Java's CDI importer cannot +read. Compatible results retain zero-copy transfer. + +The API produces a regular `PythonUDF` expression with an in-process evaluation +type. Spark's existing `ArrowEvalPython` planning rules handle aggregation, +nested calls, nondeterminism, and filter/limit pushdown. A dedicated +`InProcessArrowEvalPythonExec` extends `EvalPythonExec`, reusing its projection, +row queue, result join, and partition-evaluator path. Ordinary Python UDFs continue +to use Python workers. + +`maxRecordsPerBatch <= 0` means no row-count limit. The independent +`spark.sql.execution.arrow.maxBytesPerBatch` limit still applies when positive. +Only UDF arguments are converted to Arrow. Other columns stay in Spark rows, +buffered in a spillable queue until the results are joined back. Duplicate nested +field names in UDF arguments or declared results are rejected before Arrow Java +reads their buffers. + +Each batch uses fresh input buffers. A Python function may retain an input array; +later batches do not overwrite it. Retained arrays keep native memory alive, so +functions should release them when no longer needed. JVM input vectors and result +vectors are released on task completion, early termination and failure. + +UDF deserialization uses PySpark's bundled cloudpickle. Each task registers its +own function instance once and passes a small handle for subsequent batches. +Exception text escapes non-ASCII characters and NUL to preserve it across JEP's +JNI exception transport. +Task completion queues release of the registered function and its closure state. Imported +Python modules still share executor-wide state. Configured site-packages paths +are supplied to JEP before its first construction, so JEP itself can be found in +an archived venv. They are then processed with `site.addsitedir`, including `.pth` +files. Spark's own PySpark and Py4J distribution paths and the process `PYTHONPATH` +come first, followed by configured directories, newly discovered `.pth` paths, +and existing system paths. +Already imported modules cannot be replaced by changing the search path. + +The embedded interpreter uses isolated Python initialization and ignores Python startup +flags from the environment (including `PYTHONFAULTHANDLER` and `PYTHONDEVMODE`). Spark +explicitly restores its Python distribution paths and the executor process `PYTHONPATH` +before configured `sitePackages` paths. This also supports YARN's localized Python archives +when `SPARK_HOME` is absent. Python modules must not install process-wide signal handlers +that replace the JVM's handlers. + +Spark broadcasts, accumulators, `SparkContext.addPyFile`, and Python `TaskContext` +are not supported by this embedded runtime. Captured broadcast and accumulator +objects are rejected during serialization; functions must not access them through +imported modules either. Install modules on executors before startup, optionally +using `spark.inprocess.python.sitePackages`. Session-scoped `spark.pythonWorkerEnv.*` +settings are rejected: the shared interpreter cannot apply per-session process +environments. Configure environment variables before the executor starts, for example +with `spark.executorEnv.NAME` (or the launching environment in local mode). +In-process UDFs reject `spark.executor.pyspark.memory`, `spark.sql.pyspark.udf.profiler`, +and `spark.pythonWorkerEnv.*`. Python runs inside the JVM, so a separate Python process +memory limit cannot be applied. Use executor memory settings for sizing, and worker-based +UDFs when these Python worker features are needed. + +SQL registration through `spark.udf.register` is not supported and is rejected at registration time. +Spark Connect does not support this execution mode; both client SQL registration and +server planning reject it. The decorator accepts a `DataType` or a DDL string; DDL +strings are parsed lazily with the active Spark session. It exposes `func`, `returnType`, +`evalType`, `deterministic`, and `asNondeterministic()` along with the function's name +and docstring. +Functions must receive at least one input column (a literal also works) to determine +the batch length. Positional and keyword arguments are supported. Functions are +serialized on first use, so globals can be defined or rebound after decoration +and before that first call. The driver's Python major.minor +version must match the embedded interpreter; registration checks this before +unpickling. Python exceptions, including `SystemExit` during deserialization or +execution, are converted into task failures. Tracebacks honor the query's +`spark.sql.execution.pyspark.udf.hideTraceback.enabled`, +`spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled`, and +`spark.sql.execution.pyspark.udf.tracebackWithLocals.enabled` settings. Native process +termination remains outside this exception handling. + +## Overview + +In-process Python UDFs embed CPython directly into the Spark executor JVM using +[jep (Java Embedded Python)](https://github.com/ninia/jep), eliminating the IPC overhead of +standard Python UDFs and pandas UDFs. Data is passed to Python as +[PyArrow](https://arrow.apache.org/docs/python/) arrays via the +[Arrow C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html) — zero-copy +for compatible input and output buffers. Row-to-Arrow conversion and normalization +of sliced results still copy data. + +**Use `inprocess_udf` when:** +- You are already using `pandas_udf` for vectorized transformations and want lower latency. +- Your UDF operates on Arrow/PyArrow arrays (e.g. using `pyarrow.compute`). +- You can deploy enough executor JVMs for Python parallelism (see [Requirements](#requirements)). + +**Stick with `pandas_udf` or `udf` when:** +- You need pandas Series semantics in your UDF logic. +- You need concurrent Python invocations within a single executor. +- You are not able to install jep on executors. + +--- + +## Quick Start + +### 1. Install dependencies + +```bash +pip install "jep>=4.3.2" pyarrow cloudpickle +``` + +JEP and `org.apache.arrow:arrow-c-data` are provided dependencies and are not +bundled with Spark. Supply their JARs on the driver/executor classpaths before +starting Spark, and make the JEP native library available. Use an `arrow-c-data` +version matching Spark's Arrow Java version. Installing the Python packages alone +does not supply the Arrow Java CDI JAR. + +Building JEP from source requires a JDK, a C compiler, and development headers for +the Python version being embedded (for example, `python3.12-dev` on Ubuntu with +Python 3.12). These headers are build dependencies; running a prebuilt compatible +JEP installation does not require the development package. The corresponding +Python shared library must remain available at runtime. + +### 2. Register the plugin + +```python +spark = SparkSession.builder \ + .config("spark.plugins", + "org.apache.spark.sql.execution.python.InProcessPythonPlugin") \ + .config("spark.executor.cores", "1") \ + .config("spark.task.cpus", "1") \ + .getOrCreate() +``` + +### 3. Write and call a UDF + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import LongType + +@inprocess_udf(return_type=LongType()) +def double(x): + return pc.multiply(x, 2) + +df = spark.range(10) +df.select(double(df["id"])).show() +``` + +The function receives a `pa.Array` for each input column and must return a `pa.Array`. + +--- + +## Examples + +### String transformation + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import StringType + +@inprocess_udf(return_type=StringType()) +def upper(s): + return pc.utf8_upper(s) + +df = spark.createDataFrame([("hello",), ("world",)], ["text"]) +df.select(upper(df["text"])).show() +# +------------+ +# |upper(text) | +# +------------+ +# |HELLO | +# |WORLD | +# +------------+ +``` + +### Multi-column UDF + +A UDF receives one `pa.Array` argument per input column: + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import DoubleType + +@inprocess_udf(return_type=DoubleType()) +def weighted_sum(x, y): + return pc.add(pc.multiply(x, 0.6), pc.multiply(y, 0.4)) + +df = spark.createDataFrame([(1.0, 2.0), (3.0, 4.0)], ["x", "y"]) +df.select(weighted_sum(df["x"], df["y"])).show() +``` + +### Closure capture + +Free variables are captured by cloudpickle and frozen into the serialized UDF. The captured +value is captured at first use and shipped with the function to every executor. +Rebinding a global before the first call is reflected in the serialized function; +subsequent calls reuse the cached serialization: + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import DoubleType + +SCALE_FACTOR = 100.0 + +@inprocess_udf(return_type=DoubleType()) +def scale(x): + return pc.multiply(x, SCALE_FACTOR) +``` + +### Non-deterministic UDF + +Pass `deterministic=False` when the UDF produces different results for the same input (e.g. +random sampling). This prevents the optimizer from deduplicating or reordering calls: + +```python +import random +import pyarrow as pa +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import DoubleType + +@inprocess_udf(return_type=DoubleType(), deterministic=False) +def add_noise(x): + noise = pa.array([random.gauss(0.0, 0.01) for _ in range(len(x))]) + return pc.add(x, noise) +``` + +--- + +## Requirements + +| Requirement | Detail | +|---|---| +| Python | 3.11+; driver and embedded major.minor versions must match | +| jep | 4.3.2+ (`pip install jep`) | +| `arrow-c-data` JAR | Provided separately; match Spark's Arrow Java version | +| PyArrow | 18.0.0+ | +| cloudpickle | Bundled with PySpark | +| Python concurrency | One invocation at a time per executor (see below) | + +### Executor concurrency + +In-process UDFs use one `SharedInterpreter` on a dedicated thread per executor. +Multiple Spark tasks can share an executor, including with fractional +`spark.task.cpus`, but their Python invocations are serialized. `local[*]` therefore +works but does not provide parallel embedded Python execution. + +For throughput, consider `spark.executor.cores=1, spark.task.cpus=1` and multiple +executors. More executors also mean more JVM overhead; compare with worker-based +Arrow UDFs under the same total CPU and memory budget. + +--- + +## Deployment and Distribution + +### Local development + +For local development (e.g. `SparkSession.builder.master("local[*]")`), install jep and the +required Python packages into the virtual environment you run PySpark from. The venv's +site-packages must be supplied explicitly to the embedded interpreter. Use the +PySpark distribution from the same Spark build; a separately pip-installed PySpark +version may not contain this API or match the JVM classes. + +```bash +python3 -m venv .venv +.venv/bin/pip install "jep>=4.3.2" pyarrow +source .venv/bin/activate +``` + +JEP and Arrow CDI must be on the JVM **system classpath** before the JVM starts. +`--jars` alone only configures Spark's user classloader and is insufficient. The CDI +JAR must match the Arrow Java version in the Spark build. For example: + +```bash +JEP_DIR="$(python3 -c 'import importlib.util, pathlib; print(pathlib.Path(importlib.util.find_spec("jep").origin).parent)')" +ARROW_C_DATA_JAR=/absolute/path/to/arrow-c-data.jar +spark-submit --master 'local[1]' \ + --driver-class-path "$JEP_DIR/*:$ARROW_C_DATA_JAR" \ + --conf "spark.driver.extraLibraryPath=$JEP_DIR" \ + --conf "spark.inprocess.python.sitePackages=$(dirname "$JEP_DIR")" \ + --conf spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \ + my_app.py +``` + +### Cluster deployment — prerequisite: build and zip the venv + +Both YARN and Kubernetes support distributing a virtual environment via `--archives`. Build the +venv on a machine that matches the executor OS and Python version: + +```bash +python3 -m venv myvenv +myvenv/bin/pip install "jep>=4.3.2" pyarrow cloudpickle my-custom-lib +(cd myvenv && zip -r ../myvenv.zip .) +``` + +Adjust `python3.11` in the paths below to match the Python version in your venv. + +--- + +### YARN + +Spark extracts `--archives` to a relative path (`./myvenv/`) on each YARN container before the executor JVM +starts. The key extra config compared to local development is +`spark.executorEnv.PYSPARK_PYTHON`, which tells PySpark's Python worker to use the venv's Review Comment: **[Low, docs] `spark.executorEnv.PYSPARK_PYTHON` doesn't select the Python worker executable.** The worker executable is the function's `pythonExec`, which is decided on the driver (`self.pythonExec = os.environ.get("PYSPARK_PYTHON", "python3")` in `context.py`, or `spark.pyspark.python`), and `PythonWorkerFactory` launches exactly that; the executor environment variable is never consulted. So with this YARN example (and Kubernetes Option B at L442), ordinary `udf` / `pandas_udf` in the same application still start the node's `python3`, venv-only packages raise `ModuleNotFoundError` there, and the "consistent Python version" isn't achieved. There is also no "out-of-process fallback" for in-process UDFs. Suggestion: set `PYSPARK_PYTHON` on the submitting side or `--conf spark.pyspark.python=./myvenv/bin/python3`, as in the Python packaging guide (`python_packaging.rst`). Also, L488-490 say Spark unpacks `--archives` "at task launch time", which contradicts L351 and the code: `Executor` unpacks the initial archives in its constructor, before plugin initialization, which is what makes the relative `sitePackages` path work. ########## 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: **[Medium] `SparkFiles` and `--py-files` don't work in in-process UDFs.** Python workers call `setup_spark_files` (`worker_util.py`) before running any UDF: it sets `SparkFiles._root_directory` and `_is_running_on_worker`, and adds the root and the Python includes to `sys.path`. This bootstrap only sets up Spark's and the configured paths, and `InProcessPythonUDFBuilder` passes empty `pythonIncludes`. - After `sc.addFile("model.bin")` (or `--files` / `--archives`), `SparkFiles.get("model.bin")` in an in-process UDF fails with a bare `AssertionError` (`assert cls._sc is not None` in `SparkFiles.getRootDirectory`), while the same function works as an `arrow_udf`. Loading a model file with `SparkFiles.get` is a common pattern in the `pandas_udf` code the guide suggests migrating. - Modules shipped with `--py-files`, `spark.submit.pyFiles` or `spark.addArtifacts(..., pyfile=True)` cannot be imported (`ModuleNotFoundError` while unpickling in `_inprocess_register`), except on YARN, where they end up in the executor `PYTHONPATH`. The guide's list of unsupported features mentions broadcasts, accumulators, `SparkContext.addPyFile` and `TaskContext`, but neither of these. Suggestion: pass `SparkFiles.getRootDirectory()` from the JVM like `_spark_paths` and initialize `SparkFiles` (and `sys.path`) in the bootstrap, or document both as unsupported. ########## 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" Review Comment: **[Low] `spark.inprocess.python.sitePackages` should be a proper `ConfigEntry`.** It is a raw string read with `ctx.conf().getOption(...)` and split by hand in the plugin, so it has no `.doc`, no `.version(...)` and no validation, it is missing from `docs/configuration.md` (only the new guide describes it), and typos are silently ignored. It also starts a new top-level `spark.inprocess.*` namespace next to the existing `spark.python.*` and `spark.executor.pyspark.*` ones. Suggestion: define it in `org.apache.spark.internal.config.Python` with `ConfigBuilder(...).doc(...).version("4.4.0").stringConf.toSequence.createWithDefault(Nil)` (or `5.0.0` if this lands only in master), consider a name under an existing namespace, and read it with `ctx.conf().get(...)`. Similarly, `InProcessPythonUDFBuilder.checkConfiguration` could use `PYSPARK_EXECUTOR_MEMORY` instead of the raw `"spark.executor.pyspark.memory"` key. Note that the worker treats `0` as "no limit" (`setup_memory_limits` only applies values `> 0`), while the presence check there rejects it. -- 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]
