skip to content

Why does a Python UDF in a Spark query defeat Catalyst optimization?

level: seniorimportance: should knowfreq 52%

answer

  1. the optimizer meets a closed box
  2. the filter never reaches the file reader
  3. look for the node without a star
  4. rows leave the JVM and come back
  5. batches beat rows when Python is unavoidable

basics

~20 s

Catalyst cannot see inside a Python UDF, so it cannot push it into the data source, fold it, or compile it into generated code. The plan gains a Python-evaluation node that breaks the fused loop and ships every row into a separate Python process and back.

solid answer

~50 s

To Catalyst a Python UDF is an opaque expression: it has a signature and a return type, and nothing else the optimizer can reason about. Three consequences follow. **No pushdown** — a filter built on a UDF cannot become a `PushedFilters` entry, so the source returns every row. **No code generation** — the operator appears in the plan as `BatchEvalPython` (or `ArrowEvalPython` for a pandas UDF), which breaks the whole-stage codegen chain around it. **A process boundary** — rows are serialized out of the JVM into a Python worker, evaluated, and serialized back, per row for a plain UDF. The fix, in order of preference: express the logic with built-in `pyspark.sql.functions`, which compile into the fused loop; use a vectorized `pandas_udf` when Python is unavoidable, which moves Arrow batches instead of rows; and mark genuinely non-deterministic UDFs with `asNondeterministic()` so the optimizer does not move or re-evaluate them wrongly.

code

python · 11 lines
python
from pyspark.sql.functions import col, udf, when, regexp_extract
from pyspark.sql.types import StringType

# opaque: BatchEvalPython, no pushdown, codegen chain broken
grade_udf = udf(lambda a: "HIGH" if a > 1000 else "LOW", StringType())
slow = df.withColumn("grade", grade_udf(col("amount")))

# transparent: a Catalyst expression, fused into generated code
fast = df.withColumn(
    "grade", when(col("amount") > 1000, "HIGH").otherwise("LOW")
)

go deeper

for a junior

Know that built-in functions from pyspark.sql.functions should be your first choice and that a Python UDF is the expensive fallback, not the default way to express logic.

for a middle

Explain the three losses — pushdown, code generation, and the JVM-to-Python row boundary — and point at the BatchEvalPython node in a printed plan.

for a senior

Lead the rewrite: find the UDFs that dominate a pipeline, replace what can be expressed natively, vectorize the rest, and prove the improvement from plan diffs and runtime metrics rather than wall clock alone.

for a principal

Set the standard: UDFs need justification in review, a shared library of native expressions replaces the popular ones, and executor memory budgets account for Python worker processes outside the JVM heap.

## What Catalyst can and cannot see Catalyst optimizes by rewriting a tree of **expressions** it understands. `col("amount") > 100` is an expression with a known child, a known operator and known semantics, so the optimizer can push it below a join, hand it to a Parquet reader as a pushed filter, or generate a comparison instruction for it. A function you wrap with `udf(...)` is none of those things: it is a reference to arbitrary Python code with a declared return type. Catalyst can place it in the tree and call it, and that is all. ## The three costs **1. Optimization stops at the UDF.** No pushdown into the data source: a UDF-based predicate cannot be translated into a source filter, so every row is read and decoded. No constant folding, no simplification, no use of the predicate for partition pruning. Spark can still push *other*, transparent predicates below the UDF where semantics allow, which is why ordering your filters matters — put the cheap, transparent ones first and the UDF last. **2. Whole-stage code generation breaks.** Adjacent operators are normally fused into one compiled Java loop over rows. A Python UDF cannot be compiled into that loop, so the planner inserts a `BatchEvalPython` node, and the operators above and below it belong to different codegen stages. In a printed plan you see the `*(n)` markers stop at the Python node and resume with a new id above it. **3. Every row crosses a process boundary.** Executors run Python workers as separate OS processes. A plain Python UDF serializes each row's arguments from the JVM to the worker, runs the function, and serializes the result back. That is per-row serialization plus interpreter overhead, against a baseline of a generated loop that touches memory directly — an order of magnitude is a routine outcome. The Python workers also consume memory outside the JVM heap, which is a common cause of a container being killed for exceeding its memory limit even though the JVM heap looked fine. ## What the plan looks like ``` *(2) Project [id#10, result#44] +- BatchEvalPython [my_udf(code#11)], [result#44] +- *(1) Filter isnotnull(code#11) +- FileScan parquet [id#10,code#11] PushedFilters: [IsNotNull(code)] ``` The `BatchEvalPython` node with no `*` prefix, sitting between two codegen stages, is the signature. If the same logic had been written with built-in functions, there would be one contiguous codegen stage and possibly an extra entry in `PushedFilters`. ## The remedies, best first **Express it with built-in functions.** Most UDFs written in practice are string manipulation, conditional logic, date arithmetic or JSON extraction, and all of these exist in `pyspark.sql.functions`: `when`/`otherwise`, `regexp_extract`, `split`, `concat_ws`, `date_add`, `from_json`, `coalesce`, and higher-order functions such as `transform`, `filter` and `aggregate` for array columns. These are Catalyst expressions, so they optimize and code-generate like anything else. This rewrite is where nearly all of the win is. **Use `pandas_udf` when Python logic is genuinely required.** A pandas UDF receives an Arrow batch as a `pandas.Series` and returns one, so the boundary is crossed once per batch rather than once per row, and the data moves in Arrow's columnar format with no per-row pickling. The plan shows `ArrowEvalPython`. It remains opaque to Catalyst — no pushdown, still a codegen break — but the transport cost collapses. This is the right tool for scoring a model or applying a vectorized numeric transform. **A Scala or Java UDF** removes the process hop, since it runs inside the JVM, but stays opaque to the optimizer: still no pushdown, still no fusion into surrounding generated code. It is worth it for logic that cannot be expressed otherwise on a JVM-language team, and not worth introducing a second language for otherwise. ## Determinism, and a subtle correctness trap Spark assumes a UDF is deterministic unless you say otherwise. That assumption licenses the optimizer to reorder, duplicate or eliminate calls. A UDF that reads the clock, draws a random number or calls an external service violates it, and the symptom is results that differ between runs or between a `cache()`d and uncached DataFrame. Declare such functions with `asNondeterministic()`. A related trap: because a UDF is opaque, Spark cannot know that it fails on null input, so it will happily call it on nulls that a transparent expression would have short-circuited — handle nulls inside the function or filter them out first. ## What to say in an interview Name the plan node, the boundary and the lost rewrites, then go straight to the remedy hierarchy and the measurement: compare `explain("formatted")` before and after the rewrite, and look for the filter reappearing in `PushedFilters` and the codegen stages merging. "UDFs are slow" is a weak answer; "Catalyst cannot see through it, so pushdown and fusion are lost and rows cross into a Python process" is the one being asked for.

  • How does a pandas_udf differ from a plain Python UDF in the plan and at runtime?
    It appears as ArrowEvalPython rather than BatchEvalPython, and it receives a batch of rows as a pandas Series over Arrow instead of one row at a time. That removes per-row pickling and interpreter dispatch, often by a large factor. What it does not change is visibility: Catalyst still cannot push it into the source or fuse it into generated code.
  • Does a Scala UDF solve the problem?
    Only half of it. Running in the JVM removes the process hop and the serialization, so throughput improves substantially. But a Scala UDF is still an opaque expression to the optimizer: predicates built on it are not pushed to the data source and the surrounding operators are not fused with it. Built-in functions remain the only fully transparent option.
  • Why should a UDF that calls an external service be marked asNondeterministic()?
    Spark assumes determinism and uses that to reorder, duplicate or eliminate evaluations — it may call the function more than once per row, or evaluate it in a different position after a rewrite. For a side-effecting or time-dependent function that produces inconsistent results between runs and between cached and uncached DataFrames. Marking it restricts the rewrites the optimizer is allowed to apply.
  • You rewrote a UDF with built-in functions. How do you prove it helped?
    Diff explain("formatted") before and after: the BatchEvalPython node should be gone, the codegen stage ids should merge into one contiguous run, and any predicate involved should now appear under PushedFilters on the scan with a narrower ReadSchema. Then confirm with the SQL tab's per-node row counts that fewer rows leave the scan.

A translator can rearrange a sentence they understand, but a sealed envelope inside it has to be carried across the room, opened by someone else, and carried back — untouched and in order.

saying these in an interview costs you the question

  • Says UDFs are slow without naming pushdown or the process boundary
  • Thinks Catalyst optimizes inside a UDF once the return type is declared
  • Believes a pandas_udf becomes visible to the optimizer
  • Assumes a Scala UDF restores predicate pushdown
  • Ignores that Python workers use memory outside the JVM heap

context