Skip to content
Good engineers know 4 min read · Python 2 practice problems ↓

UDFs and pandas UDFs

When a Python function is the only way, how to write it, and why built-ins are much faster.

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

Comfortable with these? Read on.

TL;DR A UDF runs your own Python function on every row. It works, but each row is serialised, sent to a Python worker and sent back, and Catalyst cannot see inside it. Use a built-in function if one exists; if not, use a pandas UDF, which moves data in Arrow batches and runs vectorised code, typically many times faster.

Writing a UDF

PySpark · A Python 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 None yourself: 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

The executorexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more → (a JVM) holds the rows in Spark's internal binary format.

Step 2 · serialise

Rows are converted and pickled, then streamed over a socket to a separate Python worker process.

Step 3 · Python

The Python interpreter calls your function once per row, one value at a time.

Step 4 · back

Results are pickled again and sent back to the JVM, then converted to Spark's format.

Step 5 · the cost

Serialisation both ways plus per-row Python calls typically makes this several times slower than a built-in. And because CatalystCatalyst: Spark's query optimizer. It rewrites your DataFrame code into a faster equivalent plan before running it. Learn more → sees the UDF as a black box, it cannot push filters through it or prune columns it reads.

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:

PySpark · Series to Series
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"))
KindSignatureUse for
Series to Seriespd.Series → pd.SeriesRow-wise maths, string clean-up, model scoring
Series to scalarpd.Series → floatA custom aggregate in groupBy().agg()
applyInPandasgroupBy(...).applyInPandas(fn, schema)Run a pandas function on each whole group, such as fitting one model per store
mapInPandas / mapInArrowiterator of DataFrames → iteratorBatch 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.

Note: the in-browser engine does not run UDFs, so these examples are read-only.

Common mistakes

A UDF for something built in

Upper-casing, regex, date maths and conditions all exist as functions.

Wrong returnType

Spark returns null instead of raising an error.

Ignoring None

The function crashes with a Python error on the first null.

Side effects inside a UDF (API calls, counters)

It may run more than once per row. Use mapInPandas with care, or do it outside Spark.

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 questions

1. 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.

Solve: Parse UTM Parameters →

Keep going

Up next · lesson 27 of 35 · 4 min read
rangeBetween
Window frames by value instead of row count: the last 7 calendar days, not the last 7 rows.

Related lessons

Previous: JSON and structs

Primary sources: functions.udf · pandas UDFs