What you will be able to do
- Load an MLflow model as a Spark UDF and call it inside a pipeline dataset definition
- List the prerequisites for using MLflow models in a pipeline: installing MLflow and choosing the channel for Unity Catalog
- Generate a Lakeflow pipelines inference notebook from the legacy Workspace Model Registry UI
- Predict what happens to existing predictions in a streaming table when the scoring query or model changes
1.An MLflow model is just another DataFrame transformation
In Lakeflow pipelines (the product formerly called Delta Live Tables), you define each dataset as a function that returns a Spark DataFrame. Streaming inference fits in without a special serving component, because Databricks treats an MLflow model as a transformation: it takes a Spark DataFrame in and returns a Spark DataFrame out. To score a stream, you write a pipeline dataset that applies the model to incoming rows, and the pipeline handles orchestration and incremental processing. If you already have a Python script that calls an MLflow model, the docs say you can convert it into a pipeline with a few lines of code.
| Requirement | What to do |
|---|---|
| MLflow library | Run %pip install mlflow in the pipeline source; it is not installed by default |
| Imports | Import mlflow and dp (from pyspark import pipelines as dp) at the top of the source |
| Unity Catalog-enabled pipeline | Configure the pipeline to use the preview channel |
| Current channel instead | Configure the pipeline to publish to the Hive metastore |
The channel. A Unity Catalog-enabled pipeline must use the preview channel to use MLflow models. Staying on the current channel requires publishing to the Hive metastore.
Checkpoint 1 of 7· Check yourself
How does Databricks treat an MLflow model inside a pipeline?
Pipelines define datasets against DataFrames, and an MLflow model acts on a DataFrame. That is why a model can sit inside an ordinary dataset definition.
“MLflow models are treated as transformations in Databricks”Source: docs.databricks.com
Sources1
2.The three-step pattern: model URI, Spark UDF, dataset definition
The docs describe applying a model as three steps. First, get the run ID and model name of the MLflow model and combine them into a model URI. Second, use that URI to define a Spark UDF that loads the model, with mlflow.pyfunc.spark_udf. Third, call the UDF inside your table definitions, passing the feature columns as its argument.
%pip install mlflow==2.20.2
from pyspark import pipelines as dp
import mlflow
run_id= "<mlflow-run-id>"
model_name = "<the-model-name-in-run>"
model_uri = f"runs:/{run_id}/{model_name}"
loaded_model_udf = mlflow.pyfunc.spark_udf(spark, model_uri=model_uri)
@dp.materialized_view
def model_predictions():
return (
spark.read.table(<input-data>)
.withColumn("prediction", loaded_model_udf(<model-features>))
)The documented example uses @dp.materialized_view with a batch read (spark.read.table), which recomputes predictions for the whole input. The docs also say existing MLflow code can be adapted with either the @dp.table or the @dp.materialized_view decorator. @dp.table creates a streaming table, and a streaming table must read from a streaming source. To score only newly arriving rows, put the same withColumn call with the UDF inside a @dp.table function that reads a streaming source. In the longer loan-risk example in the docs, the features are a list of categorical and numeric columns passed as loaded_model_udf(struct(features)).
Checkpoint 2 of 7· Put it in order
Put the steps for using an MLflow model in a pipeline in order
- 1.Use the URI built from them to define a Spark UDF that loads the model
- 2.Call the UDF in your table definitions to produce predictions
- 3.Obtain the run ID and model name of the MLflow model
You need the run ID and model name to build the model URI. The URI is what spark_udf loads, and the resulting UDF is then called inside a dataset definition.
“Use the URI to define a Spark UDF to load the MLflow model.”Source: docs.databricks.com
Checkpoint 3 of 7· Fill the gap
Which MLflow function turns the model URI into something you can call on DataFrame columns?
run_id = "mlflow_run_id"
model_name = "the_model_name_in_run"
model_uri = f"runs:/{run_id}/{model_name}"
loaded_model_udf = mlflow.pyfunc. ? (spark, model_uri=model_uri)mlflow.pyfunc.spark_udf wraps the model as a Spark UDF that withColumn can call. load_model returns a Python model object rather than a UDF that can be used on columns.
Source: docs.databricks.comCheckpoint 4 of 7· Exam question
A model is registered in Unity Catalog and referenced by the alias `champion`. Inside a Lakeflow pipeline dataset function that reads a streaming source, which snippet correctly loads and applies that model to produce a `prediction` column?
Correct answer: B — Call `mlflow.pyfunc.spark_udf(spark, model_uri="models:/catalog.schema.model@champion")`, then apply the UDF via `.withColumn("prediction", loaded_model(struct(*feature_cols)))` on each row.
- A. `load_model` returns a plain pyfunc object whose `.predict()` expects a local pandas or numpy input, not a distributed Spark DataFrame, so it cannot be called this way on streaming data inside a table definition.
- B. This is correct: `mlflow.pyfunc.spark_udf` loads the aliased model as a Spark UDF, and wrapping the feature columns in `struct` and applying `withColumn` scores each row of the incremental DataFrame.
- C. `spark_udf` returns a callable Python UDF, not an MLlib `Transformer`, so it has no `.transform()` method; it must be applied through `withColumn` or `select` like any other UDF.
- D. This is unnecessary: `mlflow.pyfunc.spark_udf` can be called directly inside a dataset function, so hand-rewriting the model's logic as a separate `pandas_udf` discards the registered model definition for no reason.
3.Generating the inference notebook from the model registry
You don't have to write the scoring code by hand. The legacy Workspace Model Registry page describes a UI path. In the Configure model inference dialog, click the Streaming (Lakeflow pipelines) tab and choose a model version. The first two options are the current Production and Staging versions, and if you pick one, the notebook uses whichever version holds that stage when it runs. Next, browse to the input table, which in Unity Catalog workspaces you select by catalog, database and table. Finally, give the output a name.
The generated notebook builds a data transform that reads the input table and uses the MLflow PySpark inference UDF to make predictions. It stores the predictions in a live table with the name you gave. You can edit the notebook, for example to add transformations before or after the model, or to define a streaming live table as the output with schema information or data quality constraints. You can then create a new pipeline with the notebook, or add it to an existing pipeline as an extra notebook library. This feature is labelled as a preview.
Checkpoint 5 of 7· Check yourself
You select "Production" in the Model version drop-down when you generate a Lakeflow pipelines inference notebook. Later, a newer version is moved to Production. What does the notebook use on its next run?
Choosing a stage instead of a fixed version number makes the notebook resolve that stage each time it runs, so it follows the model without code changes.
“the notebook automatically uses the Production or Staging version as of the time it is run”Source: docs.databricks.com
Sources3
4.When the model changes: new rows only, unless you full-refresh
If streaming predictions are written to a streaming table, the processed-once behaviour has a consequence that is easy to miss. A row already appended to a streaming table is not queried again on later pipeline updates. If you change the query, only rows processed after the change reflect the new logic. The docs illustrate this by switching LOWER(name) to UPPER(name): existing rows stay lowercase and only new rows become uppercase. Pointing the UDF at a different model is also a change to the query, so the same rule applies to scoring logic.
They stay as they are. Existing rows keep the output of the old query, and only new rows are scored by the updated definition. To reprocess everything with the latest definition, run a full refresh (REFRESH TABLE <table_name> FULL for a streaming table created with SQL). A full refresh requeries all previous data from the source. Be careful with sources that don't keep their full history or have short retention, such as Kafka. A full refresh truncates the existing data, and old records may no longer be available to reload.
Checkpoint 6 of 7· Exam question
A fraud-detection team streams scored transactions through a Lakeflow pipeline and is concerned about duplicate predictions being written if the pipeline restarts mid-run. Assuming `transactions_bronze` is an append-only source, what should the team expect from the streaming table that applies the model?
Correct answer: A — Each transaction is processed and scored exactly once, and a restart resumes from the pipeline's checkpointed progress rather than rescoring transactions that were already processed.
- A. This is correct: a streaming table over an append-only source processes each record exactly once, and restarts resume from tracked progress instead of reprocessing rows that were already scored.
- B. Reprocessing the full history on every run describes a materialized view's batch recomputation, not a streaming table, which only handles new or changed source data incrementally.
- C. Watermarking is a separate mechanism used for time-based windowed aggregations; it is not a prerequisite for a model UDF to score rows arriving on a streaming table.
- D. There is no dedicated streaming model registration type in the Unity Catalog Model Registry; a normally registered model can be loaded with `mlflow.pyfunc.spark_udf` for both batch and streaming use.
Checkpoint 7 of 7· Check yourself
A streaming table of predictions reads from a Kafka topic with a 3-day retention period. You want last month's rows rescored with a new model. What is the risk of running a full refresh?
A full refresh reprocesses everything still in the source after truncating the table. Data that has expired from Kafka is gone from the table and cannot be recomputed.
“the full refresh truncates the existing data”Source: docs.databricks.com
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Updating the model referenced in a streaming table's definition rescores every existing row on the next pipeline update.Why is that wrong?
Rows already appended to a streaming table are not queried again. Only new rows get the new logic, and reprocessing old rows needs a full refresh.
Covered in When the model changes: new rows only, unless you full-refresh
2.MLflow is available in pipelines out of the box, so the scoring code only needs import mlflow.Why is that wrong?
Pipelines do not install MLflow by default. Install it with %pip install mlflow and import mlflow and dp at the top of the source.
Covered in An MLflow model is just another DataFrame transformation
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/ldp/transformOfficial docs
“MLflow models are treated as transformations in Databricks”
↩︎ An MLflow model is just another DataFrame transformation“To use MLflow models in a Unity Catalog-enabled pipeline, your pipeline must be configured to use the preview channel.”
↩︎ An MLflow model is just another DataFrame transformation“Call the UDF in your table definitions to use the MLflow model.”
↩︎ The three-step pattern: model URI, Spark UDF, dataset definition“you can adapt this code to a pipeline by using the @dp.table or @dp.materialized_view decorator”
↩︎ The three-step pattern: model URI, Spark UDF, dataset definition“Pipelines do not install MLflow by default”
↩︎ Exam trap 2“Pipelines do not install MLflow by default”
↩︎ Prediction“Use the URI to define a Spark UDF to load the MLflow model.”
↩︎ Checkpoint - 2.https://docs.databricks.com/aws/en/ldp/conceptsOfficial docs
“Streaming tables are always defined against streaming sources.”
↩︎ The three-step pattern: model URI, Spark UDF, dataset definition - 3.https://docs.databricks.com/aws/en/machine-learning/manage-model-lifecycle/workspace-model-registryOfficial docs
“Click the Streaming (Lakeflow pipelines) tab.”
↩︎ Generating the inference notebook from the model registry“integrates the MLflow PySpark inference UDF to perform model predictions”
↩︎ Generating the inference notebook from the model registry“define a streaming live table as output, add schema information or data quality constraints”
↩︎ Generating the inference notebook from the model registry“the notebook automatically uses the Production or Staging version as of the time it is run”
↩︎ Checkpoint - 4.
“existing rows will not update to be uppercase, but new rows will be uppercase”
↩︎ When the model changes: new rows only, unless you full-refresh“You can trigger a full refresh to requery all previous data from the source table to update all rows in the streaming table.”
↩︎ When the model changes: new rows only, unless you full-refresh“A row that has already been appended to a streaming table will not be re-queried with later updates to the pipeline.”
↩︎ Exam trap 1 - 5.https://docs.databricks.com/aws/en/sql/language-manual/sql-ref-syntax-ddl-create-streaming-tableOfficial docs
“Changes to the provided query only get reflected on new data by calling a REFRESH, not previously processed data.”
↩︎ When the model changes: new rows only, unless you full-refresh“the full refresh truncates the existing data”
↩︎ Checkpoint