PySpark courseLesson 6 of 6
PySpark course · Lesson 6 of 6
PySpark UDFs and Safer Alternatives
When Python UDFs are slow or risky in PySpark, how built-in functions and pandas UDFs compare, and what Arrow-optimised UDFs in Spark 4 change.
On this page
A user-defined function (UDF) lets you run your own Python code on each value of a column. It is sometimes necessary, but it should be the last option, not the first.
Why plain Python UDFs cost more
Built-in functions (F.upper, F.when, F.regexp_extract and hundreds more) run inside Spark’s JVM engine, where the optimiser understands them and can generate efficient code. A Python UDF is a black box:
- data must be serialised from the JVM to a Python worker process and back;
- the optimiser cannot see inside it, so it cannot push filters through it or reason about its output;
- errors and
Nonehandling are your responsibility.
Compare the plans:
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
from pyspark.sql.types import StringType
df = spark.createDataFrame([("asha",), ("ben",), (None,)], ["customer"])
upper_udf = F.udf(lambda s: s.upper() if s is not None else None, StringType())
df.select(upper_udf("customer")).explain() # an EvalPython step appears
df.select(F.upper("customer")).explain() # just a Project inside the JVM
The UDF version shows a separate Python evaluation step (BatchEvalPython, or ArrowEvalPython when Arrow is used); the built-in version is a single projection.
Option 1: use built-in functions
Most “I need a UDF” cases are covered by built-ins: string handling, dates, when/otherwise, regular expressions, arrays and maps (transform, filter, aggregate), JSON parsing. Search the pyspark.sql.functions reference before writing a UDF.
df.select(
F.when(F.col("customer").isNull(), "unknown")
.otherwise(F.initcap("customer"))
.alias("customer")
).show()
Option 2: pandas (vectorised) UDFs
If you truly need Python logic, a pandas UDF processes a whole batch as a pandas Series, using Apache Arrow for efficient transfer:
import pandas as pd
@F.pandas_udf("double")
def add_tax(amount: pd.Series) -> pd.Series:
return amount * 1.18
prices = spark.createDataFrame([(100.0,), (250.0,)], ["amount"])
prices.select(add_tax("amount").alias("with_tax")).show()
Batch processing amortises the per-row overhead and lets you use vectorised NumPy/pandas operations.
What changed in Spark 4
In recent Spark releases, regular Python UDFs can also use Arrow for data transfer (spark.sql.execution.pythonUDF.arrow.enabled, which defaults to true in Spark 4.2 when PyArrow is available). That reduces serialisation cost, but the UDF is still opaque to the optimiser and still runs in Python, so built-ins remain the better choice whenever they can express the logic.
When a UDF is reasonable
- Calling a well-tested Python library (for example a parser) that has no Spark equivalent.
- Complex business logic that would be unreadable as nested expressions, on data volumes where the cost is acceptable.
In those cases: handle None explicitly, declare the return type, keep the function pure, and unit-test it as ordinary Python.
Common mistakes
- Writing a UDF for something a built-in already does.
- Forgetting
Nonehandling, so one null crashes the job. - Declaring the wrong return type and getting silent
NULLs. - Calling external services (APIs, databases) from inside a UDF, row by row.
Interview relevance
“When should you avoid Python UDFs?” is a frequent PySpark question. See the interview answer.
Key takeaway
Prefer built-in functions, then pandas UDFs, and use plain Python UDFs only when nothing else can express the logic.
Progress is saved in this browser only. No account needed.