CertSafari
    Databricks Certified Machine Learning Associate· Lessons

    Domain 4 · Lesson 46/48

    Applying an MLflow Model in a Lakeflow Pipeline with spark_udf

    Identify how streaming inference is performed with Delta Live Tables

    10 min read
    2.08% of exam
    5 sources
    Published 2 Oct 2026
    Docs as of 30 Sep 2026

    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.

    Prerequisites before a pipeline can apply an MLflow model
    RequirementWhat to do
    MLflow libraryRun %pip install mlflow in the pipeline source; it is not installed by default
    ImportsImport mlflow and dp (from pyspark import pipelines as dp) at the top of the source
    Unity Catalog-enabled pipelineConfigure the pipeline to use the preview channel
    Current channel insteadConfigure the pipeline to publish to the Hive metastore

    Checkpoint 1 of 7· Check yourself

    How does Databricks treat an MLflow model inside a pipeline?

    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.

    Basic syntax: load the model as a Spark UDF and add a prediction column in a pipeline datasetpython
    %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. 1.Use the URI built from them to define a Spark UDF that loads the model
    2. 2.Call the UDF in your table definitions to produce predictions
    3. 3.Obtain the run ID and model name of the MLflow model

    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)

    Checkpoint 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?

    Sources12

    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?

    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?

    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?

    Sources45

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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. 1.
      “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. 2.
      “Streaming tables are always defined against streaming sources.”
      ↩︎ The three-step pattern: model URI, Spark UDF, dataset definition
    3. 3.
      “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. 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. 5.
      “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

    Spotted a mistake, or was something unclear? Tell us.