viirya commented on code in PR #58978: URL: https://github.com/apache/spark/pull/58978#discussion_r4238670221
########## docs/sql-pyspark-inprocess-udf.md: ########## @@ -0,0 +1,728 @@ +--- +layout: global +title: In-Process Python UDFs +displayTitle: In-Process Python UDFs +license: | + Licensed to the Apache Software Foundation (ASF) under one or more + contributor license agreements. See the NOTICE file distributed with + this work for additional information regarding copyright ownership. + The ASF licenses this file to You under the Apache License, Version 2.0 + (the "License"); you may not use this file except in compliance with + the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +--- + +* Table of contents +{:toc} + +## Runtime and result contract + +Each executor owns a dedicated interpreter thread. The plugin initializes the +interpreter on that thread, and task calls and shutdown are dispatched to the +same thread. The JVM is asked to allocate an 8 MiB stack for this thread; the +actual size is platform-dependent. Calls from concurrent tasks are queued on the +interpreter thread. +One task per executor is recommended for throughput, but is not a correctness requirement. +Application-level Python parallelism comes from multiple executor JVMs. +The plugin configures JEP's process-wide interpreter with hash seed `0`, matching +Spark's default Python worker seed. It must initialize before any other JEP user in +the JVM. The seed cannot change between SparkContexts in the same process; a custom +worker `PYTHONHASHSEED` does not override this embedded-runtime setting. + +Task cancellation cannot safely stop arbitrary native Python code. An interrupted +caller waits for the current invocation to finish before freeing the Arrow CDI +structures, then restores its interrupt status. A UDF that never returns can +therefore prevent its task from completing cancellation and block every subsequent +in-process UDF on that executor, including calls from other tasks, jobs, and sessions. +Recovery from a permanently hung invocation requires replacing the executor process. +Plugin shutdown stops accepting new calls and waits up to five seconds for the interpreter thread. If a call is +still running or a task still owns exported results, cleanup waits for that task to release +its CDI references; the memory remains live until cleanup completes or the process exits. Shutdown does not forcibly interrupt native +code. A new interpreter cannot start until the previous one has fully stopped. + +A scalar UDF must return a `pyarrow.Array` with exactly one element per input row. +The runtime checks the result type against the declared Spark type, including +nested fields, decimal scale, and timestamp unit. Timezone-aware timestamps are relabeled +to `spark.sql.session.timeZone` without changing their UTC instants or copying their buffers. +Timezone-naive and timezone-aware timestamps are not interchangeable. String and binary +offset widths, including nested values, are converted as needed to match +`spark.sql.execution.arrow.useLargeVarTypes`. Large, fixed-size and dictionary-encoded +representations of the declared types (`large_list`, `fixed_size_list`, `string_view`, +`binary_view`, `fixed_size_binary` and dictionary arrays) are cast to the declared type. +These conversions can allocate new buffers. Other value types must match exactly: use an +explicit PyArrow cast for numeric conversions. +Map `keys_sorted` metadata is normalized to Spark's declared map type. +Nested field nullability may differ if the actual values satisfy the declared nullability. Sliced results, including nested +child slices, are copied to remove offsets that Arrow Java's CDI importer cannot +read. Zero-length levels without a usable offsets buffer, which Arrow permits, are given +one. Compatible results retain zero-copy transfer. +Before exporting a result, the runtime performs full Arrow validation, including interior +offsets, because the JVM reads result buffers without bounds checks: a malformed result, +such as one built from raw buffers, could otherwise produce wrong values or crash the +executor. It does not validate UTF-8 in string results, because Spark strings may contain +invalid UTF-8 (for example, `CAST(X'FF' AS STRING)`). Worker-based Arrow UDFs do not +validate their results. To skip the full validation, set +`spark.sql.execution.pythonUDF.inProcess.fullValidation.enabled` to `false`; Arrow's +constant-time validation and the conversions above still apply. + +The API produces a regular `PythonUDF` expression with an in-process evaluation +type. Spark's existing `ArrowEvalPython` planning rules handle aggregation, +nested calls, nondeterminism, and filter/limit pushdown. A dedicated +`InProcessArrowEvalPythonExec` extends `EvalPythonExec`, reusing its argument extraction +and partition-evaluator path, while its evaluator buffers and joins input rows itself. +Ordinary Python UDFs continue to use Python workers. + +`maxRecordsPerBatch <= 0` means no row-count limit. The independent +`spark.sql.execution.arrow.maxBytesPerBatch` limit always applies. +Only UDF arguments are converted to Arrow. Other columns stay in Spark rows, +buffered in a spillable queue until the results are joined back. When every input +column is a UDF argument and its type, other than a decimal, an array or a map, reads back +from Arrow unchanged, the output reads those columns from the Arrow input vectors instead +of buffering the rows. +Duplicate nested field names in UDF arguments or declared results are rejected before +Arrow Java reads their buffers. + +Each batch uses fresh input buffers. A Python function may retain an input array; +later batches do not overwrite it. Retained arrays keep native memory alive, so +functions should release them when no longer needed. JVM input vectors and result +vectors are released on task completion, early termination and failure. The runtime retains +each exported result until the next invocation for that task or task cleanup, after the JVM +has released its references. The runtime drops its Python references on the interpreter +thread, so releasing JVM results does not trigger Python finalizers on Spark task threads. +Cleanup can remain queued behind another task's invocation. The rows can also be consumed +on another thread, such as a pipelined Python worker's writer. Task completion then stops +that consumer after the input row it is reading, and waits for it, but not for this +operator's Python: it releases the buffered rows at once, and the Arrow vectors when Python +returns. Reading one row can take longer when the input is another in-process UDF, whose +next row may need a batch of Python, or a blocked upstream operator. While the consumer +reads its input, task completion waits for at most one second, and then leaves the buffered +rows to the executor and deletes their spill files. For a task that otherwise succeeds, the +executor then logs "Managed memory leak detected", or fails the task if +`spark.unsafe.exceptionOnMemoryLeak` is `true`. Each row is copied before it is returned, so +it remains valid after the task releases them. + +UDF deserialization uses PySpark's bundled cloudpickle. Each task registers its +own function instance once and passes a small handle for subsequent batches. +Exception text escapes NUL, surrogates and non-BMP characters for JEP's JNI exception +transport. Other characters, including non-English BMP text, remain readable. +Task completion queues release of the registered function and its closure state. Imported +Python modules still share executor-wide state. Configured site-packages paths +are supplied to JEP before its first construction, so JEP itself can be found in +an archived venv. They are then processed with `site.addsitedir`, including `.pth` +files. Spark's own PySpark and Py4J distribution paths and the process `PYTHONPATH` +come first, followed by configured directories, newly discovered `.pth` paths, +and existing system paths. +Already imported modules cannot be replaced by changing the search path. + +The embedded interpreter uses isolated Python initialization and ignores Python startup +flags from the environment (including `PYTHONFAULTHANDLER` and `PYTHONDEVMODE`). Spark +explicitly restores its Python distribution paths and the executor process `PYTHONPATH` +before configured `sitePackages` paths. This also supports YARN's localized Python archives +when `SPARK_HOME` is absent. Python modules must not install process-wide signal handlers +that replace the JVM's handlers. JEP's automatic Java package discovery is disabled so +Java packages do not shadow Python packages. + +Executors must start with a UTF-8 locale, for example `LC_ALL=C.UTF-8` on systems that +provide it. Isolated initialization ignores `PYTHONUTF8` and `PYTHONIOENCODING` and does +not coerce an ASCII locale; the runtime warns if it detects one. Standard output and error +use line buffering and are flushed during orderly interpreter shutdown. + +Spark broadcasts, accumulators, `SparkFiles`, `--py-files`, `spark.submit.pyFiles`, +`SparkContext.addPyFile`, `SparkSession.addArtifacts(..., pyfile=True)`, and Python +`TaskContext` are not supported by this embedded runtime. Python files may incidentally +be importable on YARN through its process `PYTHONPATH`; this is not portable support for +these APIs. For archived environments, access files by their configured executor paths, +rather than `SparkFiles.get`. Captured broadcast and accumulator +objects are rejected during serialization; functions must not access them through +imported modules either. Install modules on executors before startup, optionally +using `spark.inprocess.python.sitePackages`. Session-scoped `spark.pythonWorkerEnv.*` +settings are rejected: the shared interpreter cannot apply per-session process +environments. Configure environment variables before the executor starts, for example +with `spark.executorEnv.NAME` (or the launching environment in local mode). +At query execution on the driver, in-process UDFs reject positive +`spark.executor.pyspark.memory`, `spark.sql.pyspark.udf.profiler`, +`spark.python.profile`, `spark.python.profile.memory`, and `spark.pythonWorkerEnv.*` +settings. A Python memory value of `0` means no separate limit +and is accepted. Python runs inside the JVM, so a separate Python process +memory limit cannot be applied. Use executor memory settings for sizing, and worker-based +UDFs when these Python worker features are needed. Worker-specific logging, faulthandler, +traceback-dump timers, process reuse, idle timeouts, and pipelined worker transport settings +do not apply to this mode. In particular, `spark.sql.pyspark.worker.logging.enabled` and +`spark.sql.execution.pyspark.udf.faulthandler.enabled` do not enable these worker facilities +inside the JVM. Ordinary UDFs in the same application retain their worker settings. + +SQL registration through `spark.udf.register` is not supported and is rejected at registration time. +Spark Connect does not support this execution mode; both client SQL registration and +server planning reject it, and DataFrame API calls report that Connect is unsupported. The decorator accepts a `DataType` or a DDL string; DDL +strings are parsed lazily with the active Spark session. It exposes `func`, `returnType`, +`evalType`, `deterministic`, and `asNondeterministic()` along with the function's name +and docstring. +Functions must receive at least one input column (a literal also works) to determine +the batch length. Positional and keyword arguments are supported. Functions are +serialized on first use, so globals can be defined or rebound after decoration +and before that first call. The driver's Python major.minor +version must match the embedded interpreter; registration checks this before +unpickling. Python exceptions, including `SystemExit` during deserialization or +execution, are converted into task failures. Tracebacks honor the query's +`spark.sql.execution.pyspark.udf.hideTraceback.enabled`, +`spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled`, and +`spark.sql.execution.pyspark.udf.tracebackWithLocals.enabled` settings. Native process +termination remains outside this exception handling. + +## Overview + +In-process Python UDFs embed CPython directly into the Spark executor JVM using +[jep (Java Embedded Python)](https://github.com/ninia/jep), eliminating the IPC overhead of +standard Python UDFs and pandas UDFs. Data is passed to Python as +[PyArrow](https://arrow.apache.org/docs/python/) arrays via the +[Arrow C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html) — zero-copy +for compatible input and output buffers. Row-to-Arrow conversion and normalization +of sliced results still copy data. + +**Use `inprocess_udf` when:** +- You are already using `pandas_udf` for vectorized transformations and want lower latency. +- Your UDF operates on Arrow/PyArrow arrays (e.g. using `pyarrow.compute`). +- You can deploy enough executor JVMs for Python parallelism (see [Requirements](#requirements)). + +**Stick with `pandas_udf` or `udf` when:** +- You need pandas Series semantics in your UDF logic. +- You need concurrent Python invocations within a single executor. +- You are not able to install jep on executors. + +--- + +## Quick Start + +### 1. Install dependencies + +```bash +pip install "jep>=4.3.2" pyarrow cloudpickle +``` + +JEP and `org.apache.arrow:arrow-c-data` are provided dependencies and are not +bundled with Spark. Supply their JARs on the driver/executor classpaths before +starting Spark, and make the JEP native library available. Use an `arrow-c-data` +version matching Spark's Arrow Java version. Installing the Python packages alone +does not supply the Arrow Java CDI JAR. + +Building JEP from source requires a JDK, a C compiler, and development headers for +the Python version being embedded (for example, `python3.12-dev` on Ubuntu with +Python 3.12). These headers are build dependencies; running a prebuilt compatible +JEP installation does not require the development package. The corresponding +Python shared library must remain available at runtime. + +### 2. Register the plugin + +```python +spark = SparkSession.builder \ + .config("spark.plugins", + "org.apache.spark.sql.execution.python.InProcessPythonPlugin") \ + .config("spark.executor.cores", "1") \ + .config("spark.task.cpus", "1") \ + .getOrCreate() +``` + +### 3. Write and call a UDF + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import LongType + +@inprocess_udf(return_type=LongType()) +def double(x): + return pc.multiply(x, 2) + +df = spark.range(10) +df.select(double(df["id"])).show() +``` + +The function receives a `pa.Array` for each input column and must return a `pa.Array`. + +--- + +## Examples + +### String transformation + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import StringType + +@inprocess_udf(return_type=StringType()) +def upper(s): + return pc.utf8_upper(s) + +df = spark.createDataFrame([("hello",), ("world",)], ["text"]) +df.select(upper(df["text"])).show() +# +------------+ +# |upper(text) | +# +------------+ +# |HELLO | +# |WORLD | +# +------------+ +``` + +### Multi-column UDF + +A UDF receives one `pa.Array` argument per input column: + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import DoubleType + +@inprocess_udf(return_type=DoubleType()) +def weighted_sum(x, y): + return pc.add(pc.multiply(x, 0.6), pc.multiply(y, 0.4)) + +df = spark.createDataFrame([(1.0, 2.0), (3.0, 4.0)], ["x", "y"]) +df.select(weighted_sum(df["x"], df["y"])).show() +``` + +### Closure capture + +Free variables are captured by cloudpickle and frozen into the serialized UDF. The captured +value is captured at first use and shipped with the function to every executor. +Rebinding a global before the first call is reflected in the serialized function; +subsequent calls reuse the cached serialization: + +```python +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import DoubleType + +SCALE_FACTOR = 100.0 + +@inprocess_udf(return_type=DoubleType()) +def scale(x): + return pc.multiply(x, SCALE_FACTOR) +``` + +### Non-deterministic UDF + +Pass `deterministic=False` when the UDF produces different results for the same input (e.g. +random sampling). This prevents the optimizer from deduplicating or reordering calls: + +```python +import random +import pyarrow as pa +import pyarrow.compute as pc +from pyspark.inprocess.udf import inprocess_udf +from pyspark.sql.types import DoubleType + +@inprocess_udf(return_type=DoubleType(), deterministic=False) +def add_noise(x): + noise = pa.array([random.gauss(0.0, 0.01) for _ in range(len(x))]) + return pc.add(x, noise) +``` + +--- + +## Requirements + +| Requirement | Detail | +|---|---| +| Python | 3.11+; driver and embedded major.minor versions must match | +| jep | 4.3.2+ (`pip install jep`) | +| `arrow-c-data` JAR | Provided separately; match Spark's Arrow Java version | +| PyArrow | 18.0.0+ | +| cloudpickle | Bundled with PySpark | +| Python concurrency | One invocation at a time per executor (see below) | + +### Executor concurrency + +In-process UDFs use one `SharedInterpreter` on a dedicated thread per executor. +Multiple Spark tasks can share an executor, including with fractional +`spark.task.cpus`, but their Python invocations are serialized. `local[*]` therefore +works but does not provide parallel embedded Python execution. + +For throughput, consider `spark.executor.cores=1, spark.task.cpus=1` and multiple +executors. More executors also mean more JVM overhead; compare with worker-based +Arrow UDFs under the same total CPU and memory budget. + +--- + +## Deployment and Distribution + +### Local development + +For local development (e.g. `SparkSession.builder.master("local[*]")`), install jep and the +required Python packages into the virtual environment you run PySpark from. The venv's +site-packages must be supplied explicitly to the embedded interpreter. Use the +PySpark distribution from the same Spark build; a separately pip-installed PySpark +version may not contain this API or match the JVM classes. + +```bash +python3 -m venv .venv +.venv/bin/pip install "jep>=4.3.2" pyarrow +source .venv/bin/activate +``` + +JEP and Arrow CDI must be on the JVM **system classpath** before the JVM starts. +`--jars` alone only configures Spark's user classloader and is insufficient. The CDI +JAR must match the Arrow Java version in the Spark build. For example: + +```bash +JEP_DIR="$(python3 -c 'import importlib.util, pathlib; print(pathlib.Path(importlib.util.find_spec("jep").origin).parent)')" +ARROW_C_DATA_JAR=/absolute/path/to/arrow-c-data.jar +spark-submit --master 'local[1]' \ + --driver-class-path "$JEP_DIR/*:$ARROW_C_DATA_JAR" \ + --conf "spark.driver.extraLibraryPath=$JEP_DIR" \ + --conf "spark.inprocess.python.sitePackages=$(dirname "$JEP_DIR")" \ + --conf spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \ + my_app.py +``` + +### Cluster deployment — prerequisite: build and zip the venv + +Both YARN and Kubernetes support distributing a virtual environment via `--archives`. Build the +venv on a machine that matches the executor OS and Python version: + +```bash +python3 -m venv myvenv +myvenv/bin/pip install "jep>=4.3.2" pyarrow cloudpickle my-custom-lib +(cd myvenv && zip -r ../myvenv.zip .) +``` + +Adjust `python3.11` in the paths below to match the Python version in your venv. + +--- + +### YARN + +Spark extracts `--archives` to a relative path (`./myvenv/`) on each YARN container before the executor JVM +starts. Set `spark.pyspark.python` to the venv executable if ordinary worker UDFs in the same +application should also use that environment. This does not select JEP's embedded CPython; +JEP must be built against the intended Python version. In-process UDFs do not fall back to workers. + +```bash +spark-submit \ + --master yarn \ + --deploy-mode cluster \ + --archives myvenv.zip#myvenv \ + --files /absolute/path/to/arrow-c-data.jar#arrow-c-data.jar \ + --conf 'spark.executor.extraClassPath=./myvenv/lib/python3.11/site-packages/jep/*:./arrow-c-data.jar' \ + --conf spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \ + --conf spark.executor.cores=1 \ + --conf spark.task.cpus=1 \ + --conf spark.pyspark.python=./myvenv/bin/python3 \ + --conf spark.executor.extraLibraryPath=./myvenv/lib/python3.11/site-packages/jep \ + --conf spark.inprocess.python.sitePackages=./myvenv/lib/python3.11/site-packages \ + my_app.py +``` + +For local execution on the driver, use the driver classpath and native-library +settings from the local example. A client-mode driver does not execute executor UDFs. + +--- + +### Kubernetes + +#### Option A: Custom Docker image (recommended) + +Build an executor image from the same Spark distribution used by the driver. The +image must contain Python 3.11 or newer, a matching Python shared library, and a +JDK compatible with that Spark build. Do not assume `apache/spark:latest` has these +versions or includes JEP's build dependencies. + +Install JEP and PyArrow into a known environment, such as `/opt/venv`, while +building the image. Building JEP from its source distribution additionally requires +a C compiler, the matching Python development headers, and a JDK (`JAVA_HOME` set). +Then prepare the image with: + +- the JEP JAR and matching Arrow CDI JAR in `/opt/spark/jars`; +- JEP's native library directory on the JVM library path before startup; +- `spark.inprocess.python.sitePackages` pointing to `/opt/venv`'s site-packages; +- Spark's matching `python/lib/pyspark.zip` and Py4J zip in the Spark distribution. + +The plugin adds Spark's Python distribution paths itself, ahead of any PySpark +package installed in the venv. The submission example below assumes Python 3.11 +and `/opt/venv/lib/python3.11/site-packages/jep` for the native library directory; +adjust both paths to the Python version used to build JEP. + +**Submit:** + +```bash +spark-submit \ + --master k8s://https://<k8s-api-server>:<port> \ + --deploy-mode cluster \ + --conf spark.kubernetes.container.image=my-registry/spark-inprocess:tested-build \ + --conf spark.inprocess.python.sitePackages=/opt/venv/lib/python3.11/site-packages \ + --conf spark.executor.extraJavaOptions=-Djava.library.path=/opt/venv/lib/python3.11/site-packages/jep \ + --conf spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \ + --conf spark.executor.cores=1 \ + --conf spark.task.cpus=1 \ + my_app.py +``` + +`spark.executor.extraJavaOptions=-Djava.library.path=...` locates JEP even on Kubernetes +images whose entrypoint does not propagate `spark.executor.extraLibraryPath`. +It replaces the entire JVM native library search path; append any other required native +library directories, separated by `:` on Linux. + +#### Option B: `--archives` with remote file upload + +If you cannot build a custom image, Spark on Kubernetes can distribute archives via a remote +staging area (e.g. S3 or GCS). Set `spark.kubernetes.file.upload.path` to an object storage +path that both the driver and executors can access. Provision JEP and Arrow CDI JARs +in `/opt/inprocess/jars` on every executor using an image layer or a mounted volume +before the JVM starts. The image must also have Python 3.11+ and the shared +library matching the archived JEP build. Use the same JEP version as the archived +venv. A classpath wildcard pointing inside the archive is insufficient here: the JVM expands wildcards +before Spark downloads and extracts the archive. + +```bash +spark-submit \ + --master k8s://https://<k8s-api-server>:<port> \ + --deploy-mode cluster \ + --conf spark.kubernetes.container.image=my-registry/spark-python311:tested-build \ + --conf spark.kubernetes.file.upload.path=s3a://my-bucket/spark-uploads \ + --archives myvenv.zip#myvenv \ + --conf 'spark.executor.extraClassPath=/opt/inprocess/jars/*' \ + --conf spark.plugins=org.apache.spark.sql.execution.python.InProcessPythonPlugin \ + --conf spark.executor.cores=1 \ + --conf spark.task.cpus=1 \ + --conf spark.pyspark.python=./myvenv/bin/python3 \ + --conf spark.executor.extraJavaOptions="-Djava.library.path=./myvenv/lib/python3.11/site-packages/jep" \ + --conf spark.inprocess.python.sitePackages=./myvenv/lib/python3.11/site-packages \ + my_app.py +``` + +--- + +## Configuration Reference + +### `spark.plugins` + +| Default | `(none)` | +|---|---| +| **Required value** | `org.apache.spark.sql.execution.python.InProcessPythonPlugin` | + +Registers the in-process Python plugin. This initializes the `SharedInterpreter` on each +executor at startup. Without this plugin, the driver rejects in-process UDF execution +during physical planning, including `explain()`, before submitting tasks. Missing native +dependencies are reported during plugin initialization. +Task calls and cleanup never create or restart an interpreter. + +--- + +### `spark.inprocess.python.sitePackages` + +| Default | `(none)` | +|---|---| +| **Type** | Comma-separated list of absolute or relative directory paths | + +Site-package directories supplied to JEP before interpreter construction, so the +`jep` package must be directly importable from these directories. After construction, +paths are made absolute and processed with `site.addsitedir`, including `.pth` files. +Spark distribution paths and process `PYTHONPATH` take precedence over these directories; +configured paths then take precedence over system paths for modules not yet imported. +`.pth` files are processed too late to locate JEP itself during construction. + +**When you need this:** When you distribute a Python virtual environment via `--archives` and +need packages from that venv to be importable inside UDFs. The problem is that the jep +interpreter starts with the *system* Python's `sys.path`, which does not include the distributed +venv's site-packages. Setting this config tells the plugin where to find the venv's packages. + +**Typical usage with `--archives`:** + +``` +spark.inprocess.python.sitePackages = ./myvenv/lib/python3.11/site-packages +``` + +The relative path `./myvenv/` resolves to the directory where Spark extracted your archive on +the executor node. Initial archives are localized or unpacked before executor plugin +initialization, so their site-packages directories are available when JEP starts. + +Paths cannot contain a single quote, newline, NUL, surrogate character (including +supplementary Unicode characters), comma, or the platform path separator (`:` on Linux/macOS). +Backslashes are supported. Commas separate configuration entries. + +The configured directories are fixed for the JVM lifetime once JEP has been configured, +including when a subsequent Python bootstrap fails. Stopping a SparkContext does not clear +Python's cached modules or previous `.pth` entries. Restart the executor process (or the +local driver process) before changing this configuration or replacing its installed packages. + +**Multiple paths** (comma-separated): + +``` +spark.inprocess.python.sitePackages = ./venv/lib/python3.11/site-packages,/opt/custom/lib +``` + +**When you do NOT need this:** +- Executors where all required packages are pre-installed on the system Python path. + +--- + +### Native library paths + +jep requires its native library (`libjep.so` on Linux, `libjep.dylib` on macOS) to be on the +JVM's native library path. **This must be set before the JVM starts** — `System.setProperty()` +has no effect after JVM startup, so runtime configuration is not possible. + +For YARN and Standalone, use `spark.executor.extraLibraryPath` to prepend JEP's directory +to the native library search path while preserving existing directories: + +``` +spark.executor.extraLibraryPath = ./myvenv/lib/python3.11/site-packages/jep +``` + +For Kubernetes images whose entrypoint does not propagate this setting, set +`-Djava.library.path` via `spark.executor.extraJavaOptions`. This **replaces** the JVM's +entire default native library search path. Include every other required native library +directory as well, such as Hadoop native libraries or compression codecs; otherwise those +libraries may become unavailable. Separate directories with the platform path separator. + +When using `--archives`, Spark extracts the archive to a predictable relative path (`./myvenv/`), +so the path above is stable across executor nodes without any per-node configuration. + +--- + +### `spark.executor.cores` and `spark.task.cpus` + +Multiple tasks may share an executor, including fractional `spark.task.cpus` values. +Python invocations run one at a time per executor. For throughput, consider multiple +executors with: + +``` +spark.executor.cores = 1 +spark.task.cpus = 1 +``` + +--- + +## Supported Types + +Common supported Spark SQL input/output types include: + +| Category | Types | +|---|---| +| Numeric | `ByteType`, `ShortType`, `IntegerType`, `LongType`, `FloatType`, `DoubleType` | +| Boolean | `BooleanType` | +| String / Binary | `StringType`, `BinaryType` | +| Temporal | `DateType`, `TimestampType` | +| Complex | `ArrayType`, `StructType`, `MapType` | + +Nested values must satisfy the declared nullability. Map keys cannot be null. +Only types representable by Spark's Arrow conversion and JVM Arrow accessors are +supported; this is not a guarantee for every Spark SQL type. Unsupported return +types are rejected on the driver before the function is serialized or tasks start. Review Comment: Added in a85b320: "CHAR and VARCHAR return types, including nested ones, are rejected as for other Python UDFs; declare `StringType` instead." -- 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]
