UDFs and pandas UDFs
When a Python function is the only way, how to write it, and why built-ins are much faster.
On this page
You will learn
- How to write a Python UDF and register it for SQL
- Why Python UDFs are slow, step by step
- How pandas UDFs and Arrow make Python much faster
- How to avoid UDFs with built-in functions
Read first
- select · 4 min read
- Driver, executors and the cluster manager · 4 min read
Comfortable with these? Read on.
Writing a UDF
from pyspark.sql.types import StringType @F.udf(returnType=StringType()) def mask_email(email): if email is None: return None user, _, domain = email.partition("@") return user[:1] + "***@" + domain result = users.select("id", mask_email("email").alias("masked")) # Use it from SQL too spark.udf.register("mask_email", mask_email) spark.sql("SELECT id, mask_email(email) AS masked FROM users")
- Declare the return type. If the function returns something else, Spark gives null, not an error.
- Handle
Noneyourself: the UDF receives nulls as None. - Spark may call a UDF more or fewer times than you expect (it can be re-run, or evaluated before a filter). Keep it deterministic and free of side effects, or mark it
.asNondeterministic().
Why it is slow
Step 1 · JVM
Step 2 · serialise
Step 3 · Python
Step 4 · back
Step 5 · the cost
pandas UDFs
A pandas UDF (vectorised UDF) receives whole batches as pandas Series, transferred with Apache Arrow, a columnar format both the JVM and Python read without pickling. Your code then uses vectorised pandas or NumPy operations:
import pandas as pd from pyspark.sql.functions import pandas_udf @pandas_udf("double") def with_tax(amount: pd.Series, rate: pd.Series) -> pd.Series: return (amount * (1 + rate)).round(2) result = orders.select("id", with_tax("amount", "tax_rate").alias("gross"))
| Kind | Signature | Use for |
|---|---|---|
| Series to Series | pd.Series → pd.Series | Row-wise maths, string clean-up, model scoring |
| Series to scalar | pd.Series → float | A custom aggregate in groupBy().agg() |
| applyInPandas | groupBy(...).applyInPandas(fn, schema) | Run a pandas function on each whole group, such as fitting one model per store |
| mapInPandas / mapInArrow | iterator of DataFrames → iterator | Batch processing with setup cost, such as loading a model once per partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more → |
Avoiding UDFs
Before writing a UDF, check the built-ins: the mask above is F.concat(F.substring("email", 1, 1), F.lit("***@"), F.split("email", "@")[1]). Conditions are when, parsing is regexp_extract, element-wise array work is transform, and lookups are joins. Spark 3.5+ also adds Arrow-optimised regular UDFs (useArrow=True), which cut the serialisation cost when a plain UDF is unavoidable.
Common mistakes
A UDF for something built in
Wrong returnType
Ignoring None
Side effects inside a UDF (API calls, counters)
Key takeaways
- UDFs run Python per row in a separate process.
- Serialisation and the black-box effect make them slow.
- pandas UDFs use Arrow batches and vectorised code.
- Check for a built-in function first.
Check yourself
3 questions1. A UDF declared to return IntegerType returns a Python string. What happens?
Show the answer
null. A mismatched return value becomes null silently.
2. What makes pandas UDFs faster than regular UDFs?
Show the answer
Arrow batches and vectorised operations instead of per-row pickling. Data moves in columnar batches and Python works on whole Series.
3. Why can a filter not be pushed through a UDF?
Show the answer
Catalyst cannot see what the function does. The UDF is a black box to the optimiser.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 27 of 35 · 4 min readrangeBetween
Window frames by value instead of row count: the last 7 calendar days, not the last 7 rows.
Related lessons
String functionsPySpark functions · 3 min read
Array functionsSpark internals · 3 min read
Python Data Source API
Previous: JSON and structs
Primary sources: functions.udf · pandas UDFs