What you will be able to do
- Explain what makes a pandas UDF vectorized and why it is faster than a row-at-a-time Python UDF
- Create a pandas UDF with pandas_udf as a decorator or as an explicit function wrapper, choosing a returnType
- Write a Series to Series pandas UDF that respects the batch and same-length rules
- Invoke a pandas UDF from select, withColumn and, after registration, from Spark SQL
Key concept
Vectorized (pandas) UDF — A user-defined function that Spark feeds whole batches of rows, transferred with Apache Arrow, as pandas objects. Your code works on a pandas Series per batch, not on one Python value per row.
1.What a pandas UDF is and how Spark runs it
An ordinary Python UDF works row-at-a-time: Spark hands your function one value, gets one value back, and repeats this for every row. A pandas UDF, also called a vectorized UDF, does two things differently. It moves data between the JVM and Python with Apache Arrow, and it gives your function pandas objects to work on. Because of this, one call to your function processes a whole batch of rows with pandas' vectorized operations. Spark splits the column into batches, calls your function once per batch, and concatenates the results back into a column. Databricks says this design can increase performance up to 100x compared to row-at-a-time Python UDFs. That speed is the reason to choose a pandas UDF when built-in Spark functions can't express the logic you need.
Checkpoint 1 of 7· Check yourself
Which pair of technologies does a pandas UDF combine?
Arrow moves the batches between Spark and the Python worker in a columnar format, and pandas is the API your function uses on each batch.
“uses Apache Arrow to transfer data and pandas to work with the data”Source: docs.databricks.com
Sources1
2.Creating one: decorator or explicit wrapper
You create every pandas UDF with pandas_udf from pyspark.sql.functions, and you can use it in two ways. The first is to put @pandas_udf(...) as a decorator above the function definition, so the name you define is the UDF. The second is to write a plain Python function and then call pandas_udf(f, returnType=...) on it. That gives you two names: the original function, which still runs on local pandas data, and the new UDF object. Either way, no additional configuration is required.
# Declare the function and create the UDF
def multiply_func(a: pd.Series, b: pd.Series) -> pd.Series:
return a * b
multiply = pandas_udf(multiply_func, returnType=LongType())@pandas_udf("string")
def to_upper(s: pd.Series) -> pd.Series:
return s.str.upper()pandas_udf takes three parameters. returnType can be a DataType object such as LongType() or a DDL-formatted string such as "long" or "string". functionType is left over from before Spark 3.0, when you picked the UDF kind with a PandasUDFType enum. Now Spark works out the kind from your Python type hints (for example pd.Series -> pd.Series), and the reference docs prefer type hints over functionType.
| Parameter | Type | What it does |
|---|---|---|
| f | function | Optional. The Python function to wrap, when pandas_udf is used as a standalone wrapper rather than a decorator |
| returnType | pyspark.sql.types.DataType or str | Optional. The UDF's return type, as a DataType object or a DDL-formatted type string |
| functionType | int | Optional. A PandasUDFType enum value, default SCALAR, kept for compatibility. Type hints are encouraged instead |
Checkpoint 2 of 7· Fill the gap
This sample should turn multiply_func into a vectorized UDF. Which function completes it?
def multiply_func(a: pd.Series, b: pd.Series) -> pd.Series:
return a * b
multiply = ? (multiply_func, returnType=LongType())pandas_udf wraps the function as an Arrow-based vectorized UDF. Plain udf would make a row-at-a-time UDF, applyInPandas is a grouped function API, and register exposes a function to SQL.
Source: docs.databricks.comCheckpoint 3 of 7· Check yourself
A colleague writes @pandas_udf(IntegerType(), PandasUDFType.SCALAR) with no type hints. What do the docs recommend instead?
functionType only exists for compatibility. Spark works out the UDF kind from the type hints, and the docs encourage using them.
“Default: SCALAR. This parameter exists for compatibility. Using Python type hints is encouraged.”Source: docs.databricks.com
Checkpoint 4 of 7· Exam question
A data engineer wants to add a `discounted_price` column to a large Spark DataFrame `orders` by multiplying every value in the `price` column by 0.9, using a vectorized pandas UDF that operates on the whole column at once via Apache Arrow rather than row by row. Which implementation correctly achieves this?
Correct answer: A — Decorate a function typed `s: pd.Series -> pd.Series` with `@pandas_udf(DoubleType())`, then call it inside `withColumn("discounted_price", apply_discount(orders.price))` so Arrow transfers the whole column in batches.
- A. This is correct because pandas_udf with a Series-to-Series signature receives the full price column as a pandas Series through Arrow, multiplies it in one vectorized operation, and returns a same-length Series that withColumn can attach directly to the DataFrame.
- B. The plain udf decorator does not use Arrow batching, so even though the function body is identical, Spark serializes and deserializes each row individually with Python pickle, losing the vectorized performance benefit pandas_udf is meant to provide.
- C. Returning a scalar float from `s.mean() * 0.9` makes this a grouped-aggregate pandas UDF, which is meant to be called inside groupby().agg() to produce one value per group, not to multiply every row's own price and attach it as a new column.
- D. pandas_udf-wrapped functions are called like any other Spark column function, so passing the DataFrame and a column name string as two arguments produces a TypeError instead of the expected Column object needed inside withColumn.
- E. Combining the deprecated PandasUDFType.SCALAR argument with modern type-hint-based pandas_udf syntax mixes two incompatible calling conventions; PandasUDFType has been removed in current Spark releases in favor of inferring the UDF type entirely from the function's type hints, so this call raises an error.
3.The Series to Series contract
The most common kind is Series to Series, which you use to vectorize a scalar, row-wise computation. Its type hint is pd.Series, ... -> pd.Series: each input column arrives as a pandas Series, and the function must return a Series of the same length. Remember that "length" means the length of the current batch, not of the whole column. Spark calls the function once per batch, so each input Series holds only that batch's rows. If you wrote s - s.mean(), you would subtract each batch's own mean, which is not the column mean you probably wanted.
Checkpoint 5 of 7· Check yourself
Inside a Series to Series pandas UDF, what does len(s) return for an input Series s?
Spark splits the column into batches and calls the function once per batch, so each Series contains just one batch.
“is the length of an internal batch used for each call to the function”Source: spark.apache.org
The function is still plain pandas code, which helps with testing. With the explicit-wrapper style, multiply_func is still an ordinary Python function, so you can call it on a local pd.Series and check the output before Spark is involved. Only the wrapped multiply object is a Spark UDF.
# The function for a pandas_udf should be able to execute with local pandas data
x = pd.Series([1, 2, 3])
print(multiply_func(x, x))
# 0 1
# 1 4
# 2 9
# dtype: int64Sources1
4.Invoking: DataFrame APIs and Spark SQL
A pandas UDF generally behaves like a regular PySpark function. You call it with columns, either as col(...) objects or as column-name strings, and the result is a column expression. You can use that expression in select or withColumn. By default, the output column is named after the function and its arguments, for example multiply_func(x, x) or to_upper(name). A Series to Series UDF can also take keyword arguments, and the generated column name records each binding:
@pandas_udf(returnType=IntegerType())
def calc(a: pd.Series, b: pd.Series) -> pd.Series:
return a + 10 * b
spark.range(2).select(calc(b=sf.col("id") * 10, a=sf.col("id"))).show()To call the UDF from SQL, register it with spark.udf.register(name, f). That method accepts a plain Python function, a row-at-a-time udf, or a pandas_udf. One detail matters here: register's own returnType argument applies only to a plain Python function. A pandas UDF already has its return type, so you pass just the name and the UDF.
@pandas_udf("integer")
def add_one(s: pd.Series) -> pd.Series:
return s + 1
spark.udf.register("add_one", add_one)
spark.sql("SELECT add_one(id) FROM range(3)").collect()
# [Row(add_one(id)=1), Row(add_one(id)=2), Row(add_one(id)=3)]Checkpoint 6 of 7· Check yourself
You register an existing pandas UDF with spark.udf.register("add_one", add_one, LongType()). What is wrong with this call?
register's returnType applies only when f is a plain Python function. A pandas UDF already carries its return type from pandas_udf.
“Only valid when f is a plain Python function, not when f is already a user-defined function.”Source: docs.databricks.com
Checkpoint 7 of 7· Exam question
A data scientist needs to apply a machine learning model that is expensive to load to every row of a `features` column across a very large DataFrame. To avoid reloading the model once per row, which pandas UDF signature should be used so the model can be loaded once per Python worker and reused across many batches?
Correct answer: A — Type the function as `Iterator[pd.Series] -> Iterator[pd.Series]`, load the model once before the loop over the iterator, and `yield` predictions for each incoming batch of the `features` column inside the loop.
- A. The iterator-to-iterator signature runs the function once per worker process, so loading the model before entering the for-loop over the iterator means the model is initialized a single time and then reused to score every batch of the features column that arrives, avoiding repeated load costs.
- B. The scalar Series-to-Series signature is invoked once per batch by Spark, so loading the model inside the function body still reloads it for every batch rather than once per worker, which does not solve the expensive load problem described.
- C. A Series-to-scalar signature is a grouped-aggregate pandas UDF meant to collapse each group into one value inside groupby().agg(); it does not apply the model row-by-row across the whole features column the way this scenario requires.
- D. Using the iterator signature but reloading the model inside the loop still recreates the model on every batch instead of once per worker, defeating the purpose of choosing the iterator-based signature in the first place.
- E. PandasUDFType has been removed in current Spark releases, so pairing it with modern type-hint-based pandas_udf syntax is not a supported calling convention and this definition raises an error rather than running successfully.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Each call to a Series to Series pandas UDF receives the entire column, so whole-column statistics like s.mean() are safe to compute inside it.Why is that wrong?
Spark calls the function once per batch of rows and concatenates the results. Any statistic computed inside the function covers only the current batch.
Covered in The Series to Series contract
2.When registering a pandas UDF for SQL, you must pass its return type again as the third argument to spark.udf.register.Why is that wrong?
register's returnType applies only to plain Python functions. A pandas UDF already has its return type and is registered with just a name and the UDF.
Covered in Invoking: DataFrame APIs and Spark SQL
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/udf/pandasOfficial docs
“pandas UDFs allow vectorized operations that can increase performance up to 100x compared to row-at-a-time Python UDFs.”
↩︎ What a pandas UDF is and how Spark runs it“Spark runs a pandas UDF by splitting the data into batches of rows, calling the function for each batch, and then concatenating the results.”
↩︎ What a pandas UDF is and how Spark runs it“The Python function must accept a pandas Series as an input and return a pandas Series of the same length.”
↩︎ The Series to Series contract“You use a Series to Series pandas UDF to vectorize scalar operations.”
↩︎ The Series to Series contract“You can use them with APIs such as select and withColumn.”
↩︎ Invoking: DataFrame APIs and Spark SQL“A pandas user-defined function (UDF)—also known as vectorized UDF—is a user-defined function that uses Apache Arrow to transfer data”
↩︎ Key concept“Spark runs a pandas UDF by splitting the data into batches of rows, calling the function for each batch, and then concatenating the results.”
↩︎ Exam trap 1“uses Apache Arrow to transfer data and pandas to work with the data”
↩︎ Checkpoint - 2.
“A Pandas UDF is defined using the pandas_udf as a decorator or to wrap the function, and no additional configuration is required.”
↩︎ Creating one: decorator or explicit wrapper“The return type of the user-defined function. The value can be either a DataType object or a DDL-formatted type string.”
↩︎ Creating one: decorator or explicit wrapper“A Pandas UDF behaves as a regular PySpark function API in general.”
↩︎ Invoking: DataFrame APIs and Spark SQL“Default: SCALAR. This parameter exists for compatibility. Using Python type hints is encouraged.”
↩︎ Checkpoint - 3.https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.functions.pandas_udf.htmlSecondary source
“Prior to Spark 3.0, the pandas UDF used functionType to decide the execution type”
↩︎ Creating one: decorator or explicit wrapper“is the length of an internal batch used for each call to the function”
↩︎ Checkpoint - 4.
“The user-defined function can be either row-at-a-time or vectorized.”
↩︎ Invoking: DataFrame APIs and Spark SQL“Only valid when f is a plain Python function, not when f is already a user-defined function.”
↩︎ Exam trap 2“Only valid when f is a plain Python function, not when f is already a user-defined function.”
↩︎ Checkpoint