What you will be able to do
- Write a Series to Series pandas UDF and explain how Spark batches its input
- Use an Iterator pandas UDF to load a model once and reuse it for every batch
- Choose between mapInPandas and groupBy().applyInPandas for DataFrame-level scoring, and explain the memory risk of each
1.Series to Series: the basic vectorized scoring function
You can also write the scoring step yourself instead of using mlflow.pyfunc.spark_udf. The building block is a pandas UDF: a Python function with type hints, wrapped by pandas_udf, either as a decorator or as a direct call. The simplest kind is Series to Series. The function takes one or more pandas Series and returns a single Series of the same length. You apply it with select or withColumn. The next section shows the Iterator variant, which is the form to use when the function must load a model.
import pandas as pd
from pyspark.sql.functions import col, pandas_udf
from pyspark.sql.types import LongType
# 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())
# 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))The comment in that sample is the habit that matters for inference: the plain Python function must work on local pandas data. You can check your prediction logic on a few rows in the notebook before Spark spreads it over the table. Two rules come with this UDF type. The output must be exactly as long as the input batch. The returned values must also match the declared returnType, because PySpark warns that converting a mismatched type is "not guaranteed to be correct".
Checkpoint 1 of 4· Check yourself
A Series to Series scoring UDF drops rows it can't score, so it returns fewer values than it receives. What is wrong?
A Series to Series pandas UDF must return one output value per input row, so a shorter Series is invalid. If you need output of a different length, use a function API such as mapInPandas.
“The Python function must accept a pandas Series as an input and return a pandas Series of the same length.”Source: docs.databricks.com
2.Iterator UDFs: load the model once, score every batch
A Series to Series function is called once per batch. Any setup code inside it runs again for every batch. For inference, the costly setup is loading the model. The Iterator of Series to Iterator of Series UDF fixes this. Its function receives an iterator of batches, not a single batch, so you can do the setup once and then loop. The Databricks docs give model loading as the example: the UDF is useful when execution requires "initializing some state, for example, loading a machine learning model file to apply inference to every input batch".
The Iterator type hint comes from Python's typing module, so import it before defining the UDF:
import pandas as pd
from typing import Iterator
from pyspark.sql.functions import col, pandas_udf, struct# In the UDF, you can initialize some state before processing batches.
# Wrap your code with try/finally or use context managers to ensure
# the release of resources at the end.
y = 1 # value captured by the UDF closure
@pandas_udf("long")
def plus_y(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
try:
for x in batch_iter:
yield x + y
finally:
pass # release resources here, if anyFor a Python MLflow model, the docs show how to load it as a generic Python function and score data with it:
model = mlflow.pyfunc.load_model(model_path)
model.predict(model_input)Combining the two snippets gives the scoring pattern: put the mlflow.pyfunc.load_model(model_path) call where y = 1 is, before the loop, and in the loop yield model.predict(x) where the sample yields x + y. The model is then loaded once and predict runs on each batch. (This combination is our reading of the two documented pieces, not a single doc example.) The model_path can take several forms, such as a run-relative path (runs:/{run_id}/{model-path}) or a registered model path (models:/{model_name}/{model_stage}). The length rule moves up a level: the total output across all yielded batches must equal the total input. The wrapped UDF takes one Spark column. When the model needs several feature columns, use the Iterator of multiple Series variant. Its function receives an iterator of tuples of Series, and its type hint is Iterator[Tuple[pandas.Series, ...]] -> Iterator[pandas.Series].
Checkpoint 2 of 4· Fill the gap
Which type hint makes this an Iterator pandas UDF that can set up state once?
@pandas_udf("long")
def plus_one(batch_iter: ? [pd.Series]) -> Iterator[pd.Series]:
for x in batch_iter:
yield x + 1With an Iterator[pd.Series] input hint, Spark passes a stream of batches to the function. Code that runs before the loop, such as loading a model, then runs once rather than once per batch.
Checkpoint 3 of 4· Exam question
A team registered an MLflow model in Unity Catalog under the three-level name `retail.forecasting.demand_model` and wants to load version 3 of that model to run pandas-based batch predictions in a notebook. Which model loading call is correct?
Correct answer: A — `mlflow.pyfunc.load_model("models:/retail.forecasting.demand_model/3")` — the three-level name plus the numeric version identifies the exact Unity Catalog model version to load.
- A. A Unity Catalog registered model URI uses the `models:/` scheme followed by the full three-level `catalog.schema.model` name and a version or alias, so this string correctly identifies version 3 of the named model for loading.
- B. The catalog and schema segments are required parts of a Unity Catalog registered model's identity and are never inferred from the current session, so dropping them to just the leaf name would fail to resolve the model.
- C. The `models:/` scheme with slash-separated segments is required, not a bare colon-separated three-level name; this string is not a valid MLflow model URI regardless of registry backend.
- D. Unity Catalog registered models are addressed through the `models:/` URI scheme, not a fixed DBFS filesystem path, so this string does not point to a valid registered model location.
3.DataFrame-level scoring with pandas function APIs
pandas UDFs work on columns. The pandas function APIs work on whole DataFrames: you apply a native Python function that takes and returns pandas objects directly to a PySpark DataFrame. They share the internal machinery of pandas UDFs, including Arrow and the same configurations, but Python type hints are optional. Two of them matter for scoring.
DataFrame.mapInPandas() passes an iterator of pandas DataFrames to your function and expects an iterator of pandas DataFrames back. Unlike a Series to Series UDF, it can return output of any length. Because the function loops over an iterator, it has the same shape as an Iterator UDF, so code placed before the loop would run once per function call. The docs don't show a model-scoring example for it.
df = spark.createDataFrame([(1, 21), (2, 30)], ("id", "age"))
def filter_func(iterator):
for pdf in iterator:
yield pdf[pdf.id == 1]
df.mapInPandas(filter_func, schema=df.schema).show()groupBy().applyInPandas() follows the split-apply-combine pattern. It splits the data with groupBy, runs your function on each group as a pandas DataFrame, and combines the results. You must supply the function and an output schema. This fits per-group logic. The cost is memory. A group isn't streamed in batches: the whole group is loaded into memory before the function runs, and maxRecordsPerBatch doesn't apply to groups. Skewed group sizes can therefore cause out-of-memory errors. The cogrouped variant, groupby().cogroup().applyInPandas(), has the same restriction.
| Construct | Function input to output | Output length | Fit for inference |
|---|---|---|---|
| Series to Series pandas UDF | pd.Series to pd.Series, one batch per call | Same as the input batch | Simple vectorized scoring of a column |
| Iterator of Series pandas UDF | Iterator[pd.Series] to Iterator[pd.Series] | Total output equals total input | Load the model once, then score every batch |
| mapInPandas | Iterator of pandas.DataFrame to iterator of pandas.DataFrame | Arbitrary | Output may be any length |
| groupBy().applyInPandas | pandas.DataFrame per group to pandas.DataFrame | Defined by the output schema | Per-group logic; the full group must fit in memory |
Checkpoint 4 of 4· Match them up
Match each construct to how it behaves
Tap a term, then the definition that fits it.
The constructs differ in what they receive (one batch, a stream of batches, or a whole group) and in how long their output can be. Both decide whether a model can be loaded once and whether memory stays bounded.
“It can return output of arbitrary length in contrast to some pandas UDFs such as Series to Series.”Source: docs.databricks.com
Sources4
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.A Series to Series pandas UDF is the best choice when the model is expensive to load, because it is the simplest type.Why is that wrong?
A Series to Series function runs once per batch, so setup inside it repeats every time. The Iterator UDF exists so that state, such as a loaded model file, is set up once and reused for every input batch.
Covered in Iterator UDFs: load the model once, score every batch
2.groupBy().applyInPandas streams each group in Arrow batches, so large groups are safe.Why is that wrong?
Each group is loaded into memory in full before the function runs, and maxRecordsPerBatch doesn't apply to groups. Skewed groups can cause out-of-memory errors.
Covered in DataFrame-level scoring with pandas function APIs
Practise it for real
Write a vectorized pandas UDF, check it on local pandas data, run it on a Spark DataFrame, then move to the Iterator form where a model is loaded once and applied to each batch.
1.In a Databricks notebook, define multiply_func(a: pd.Series, b: pd.Series) -> pd.Series that returns a * b, and wrap it with pandas_udf(multiply_func, returnType=LongType()).
Why: A pandas UDF is a plain pandas function with type hints, registered with a Spark return type.
You should see: The cell runs without errors and defines both multiply_func and multiply.
2.Call multiply_func(x, x) directly on x = pd.Series([1, 2, 3]).
Why: The function must run on local pandas data. This is where you would test prediction logic on a few rows.
You should see: A pandas Series with the values 1, 4, 9.
3.Create df = spark.createDataFrame(pd.DataFrame(x, columns=["x"])) and run df.select(multiply(col("x"), col("x"))).show().
Why: Spark now splits the column into batches, calls the function for each batch, and concatenates the results.
You should see: A one-column Spark result named multiply_func(x, x) with the rows 1, 4, 9.
4.Add
from typing import Iterator, then rewrite the function as an Iterator UDF, plus_y(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series], with a try/finally around the loop, and run it with df.select.Why: The iterator form is where you would load a model once before the loop and release it in finally. Without the typing import, Iterator is undefined.
You should see: A one-column result with each input value increased by y (2, 3, 4 when y = 1).
5.If you have a logged MLflow model, replace y = 1 with model = mlflow.pyfunc.load_model(model_path) placed before the loop, and replace x + y with model.predict(x). Use a model whose input is a single column that matches df's column.
Why: This is the batch inference pattern itself: the model is loaded once, then predict runs on every batch.
You should see: A one-column result holding the model's predictions, with the same number of rows as the input.
Stuck? Get a nudge
If the Spark result doesn't match the local result, check that the returned Series has the same length as its input and a dtype that matches the declared returnType.
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
“You use a Series to Series pandas UDF to vectorize scalar operations. You can use them with APIs such as select and withColumn.”
↩︎ Series to Series: the basic vectorized scoring function“The length of the entire output in the iterator should be the same as the length of the entire input.”
↩︎ Iterator UDFs: load the model once, score every batch“The underlying Python function takes an iterator of a tuple of pandas Series.”
↩︎ Iterator UDFs: load the model once, score every batch“initializing some state, for example, loading a machine learning model file to apply inference to every input batch”
↩︎ Exam trap 1“The Python function must accept a pandas Series as an input and return a pandas Series of the same length.”
↩︎ Checkpoint - 2.https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.functions.pandas_udf.htmlSecondary source
“The conversion is not guaranteed to be correct and results should be checked for accuracy by users.”
↩︎ Series to Series: the basic vectorized scoring function“but is the length of an internal batch used for each call to the function.”
↩︎ Prediction - 3.https://docs.databricks.com/aws/en/mlflow/modelsOfficial docs
“For Python MLflow models, an additional option is to use mlflow.pyfunc.load_model() to load the model as a generic Python function.”
↩︎ Iterator UDFs: load the model once, score every batch - 4.
“Python type hints are optional in pandas function APIs.”
↩︎ DataFrame-level scoring with pandas function APIs“The configuration for maxRecordsPerBatch is not applied on groups”
↩︎ DataFrame-level scoring with pandas function APIs“The underlying function takes and outputs an iterator of pandas.DataFrame.”
↩︎ DataFrame-level scoring with pandas function APIs“All data for a group is loaded into memory before the function is applied.”
↩︎ Exam trap 2“It can return output of arbitrary length in contrast to some pandas UDFs such as Series to Series.”
↩︎ Checkpoint