# Intro to MLOps using ZenML

## 🌍 Overview

This repository is a minimalistic MLOps project intended as a starting point to learn how to put ML workflows in production. It features: 

- A feature engineering pipeline that loads data and prepares it for training.
- A training pipeline that loads the preprocessed dataset and trains a model.
- A batch inference pipeline that runs predictions on the trained model with new data.

Follow along this notebook to understand how you can use ZenML to productionalize your ML workflows!

<img src=".assets/pipeline_overview.png" width="70%" alt="Pipelines Overview">

## Run on Colab

You can use Google Colab to see ZenML in action, no signup / installation
required!

[![Open In Colab](https://colab.research.google.com/assets/colab-badge.svg)](
https://colab.research.google.com/github/zenml-io/zenml/blob/main/examples/quickstart/quickstart.ipynb)

# 👶 Step 0. Install Requirements

Let's install ZenML to get started. First we'll install the latest version of
ZenML as well as the `sklearn` integration of ZenML:

In [None]:
!pip install "zenml[server]"

In [None]:
from zenml.environment import Environment

if Environment.in_google_colab():
    # Install Cloudflare Tunnel binary
    !wget -q https://github.com/cloudflare/cloudflared/releases/latest/download/cloudflared-linux-amd64.deb && dpkg -i cloudflared-linux-amd64.deb


In [None]:
!zenml integration install sklearn mlflow -y

import IPython
IPython.Application.instance().kernel.do_shutdown(restart=True)

Please wait for the installation to complete before running subsequent cells. At
the end of the installation, the notebook kernel will automatically restart.

Optional: If you are using [ZenML Cloud](https://zenml.io/cloud), execute the following cell with your tenant URL. Otherwise ignore.

In [None]:
zenml_server_url = "PLEASE_UPDATE_ME"  # in the form "https://URL_TO_SERVER"

!zenml connect --url $zenml_server_url

In [1]:
# Initialize ZenML and set the default stack
!zenml init

!zenml stack set default

[?25l[1;35mInitializing the ZenML global configuration version to 0.52.0[0m
[1;35mCreating database tables[0m
[32m⠋[0m Initializing ZenML repository at 
/home/apenner/PycharmProjects/template-starter/template.
[2K[1A[2K[1A[2K[32m⠙[0m Initializing ZenML repository at 
/home/apenner/PycharmProjects/template-starter/template.
[2K[1A[2K[1A[2K[32m⠹[0m Initializing ZenML repository at 
/home/apenner/PycharmProjects/template-starter/template.
[2K[1A[2K[1A[2K[32m⠸[0m Initializing ZenML repository at 
/home/apenner/PycharmProjects/template-starter/template.
[2K[1A[2K[1A[2K[32m⠼[0m Initializing ZenML repository at 
/home/apenner/PycharmProjects/template-starter/template.
[2K[1A[2K[1A[2K[32m⠴[0m Initializing ZenML repository at 
/home/apenner/PycharmProjects/template-starter/template.
[1;35mCreating default workspace 'default' ...[0m
[1;35mCreating default stack in workspace default...[0m
[2K[1A[2K[1A[2K[32m⠦[0m Initializing ZenML repository at

In [2]:
# Do the imports at the top
from typing_extensions import Annotated
from sklearn.datasets import load_breast_cancer

import random
import pandas as pd
from zenml import step, ExternalArtifact, pipeline, ModelVersion, get_step_context
from zenml.client import Client
from zenml.logger import get_logger
from uuid import UUID

from typing import Optional, List

from zenml import pipeline

from steps import (
    data_loader,
    data_preprocessor,
    data_splitter,
    model_evaluator,
    inference_preprocessor
)

from zenml.logger import get_logger

logger = get_logger(__name__)

# Initialize the ZenML client to fetch objects from the ZenML Server
client = Client()

## 🥇 Step 1: Load your data and execute feature engineering

We'll start off by importing our data. In this quickstart we'll be working with
[the Breast Cancer](https://archive.ics.uci.edu/dataset/17/breast+cancer+wisconsin+diagnostic) dataset
which is publicly available on the UCI Machine Learning Repository. The task is a classification
problem, to predict whether a patient is diagnosed with breast cancer or not.

When you're getting started with a machine learning problem you'll want to do
something similar to this: import your data and get it in the right shape for
your training. ZenML mostly gets out of your way when you're writing your Python
code, as you'll see from the following cell.

<img src=".assets/feature_engineering_pipeline.png" width="30%" alt="Feature engineering pipeline" />

In [3]:
@step
def data_loader_simplified(
    random_state: int, is_inference: bool = False, target: str = "target"
) -> Annotated[pd.DataFrame, "dataset"]:  # We name the dataset 
    """Dataset reader step."""
    dataset = load_breast_cancer(as_frame=True)
    inference_size = int(len(dataset.target) * 0.05)
    dataset: pd.DataFrame = dataset.frame
    inference_subset = dataset.sample(inference_size, random_state=random_state)
    if is_inference:
        dataset = inference_subset
        dataset.drop(columns=target, inplace=True)
    else:
        dataset.drop(inference_subset.index, inplace=True)
    dataset.reset_index(drop=True, inplace=True)
    logger.info(f"Dataset with {len(dataset)} records loaded!")
    return dataset


The whole function is decorated with the `@step` decorator, which
tells ZenML to track this function as a step in the pipeline. This means that
ZenML will automatically version, track, and cache the data that is produced by
this function as an `artifact`. This is a very powerful feature, as it means that you can
reproduce your data at any point in the future, even if the original data source
changes or disappears. 

Note the use of the `typing` module's `Annotated` type hint in the output of the
step. We're using this to give a name to the output of the step, which will make
it possible to access it via a keyword later on.

You'll also notice that we have included type hints for the outputs
to the function. These are not only useful for anyone reading your code, but
help ZenML process your data in a way appropriate to the specific data types.

ZenML is built in a way that allows you to experiment with your data and build
your pipelines as you work, so if you want to call this function to see how it
works, you can just call it directly. Here we take a look at the first few rows
of your training dataset.

In [4]:
df = data_loader_simplified(random_state=42)
df.head()

[1;35mDataset with 541 records loaded![0m


Unnamed: 0,mean radius,mean texture,mean perimeter,mean area,mean smoothness,mean compactness,mean concavity,mean concave points,mean symmetry,mean fractal dimension,...,worst texture,worst perimeter,worst area,worst smoothness,worst compactness,worst concavity,worst concave points,worst symmetry,worst fractal dimension,target
0,17.99,10.38,122.8,1001.0,0.1184,0.2776,0.3001,0.1471,0.2419,0.07871,...,17.33,184.6,2019.0,0.1622,0.6656,0.7119,0.2654,0.4601,0.1189,0
1,20.57,17.77,132.9,1326.0,0.08474,0.07864,0.0869,0.07017,0.1812,0.05667,...,23.41,158.8,1956.0,0.1238,0.1866,0.2416,0.186,0.275,0.08902,0
2,19.69,21.25,130.0,1203.0,0.1096,0.1599,0.1974,0.1279,0.2069,0.05999,...,25.53,152.5,1709.0,0.1444,0.4245,0.4504,0.243,0.3613,0.08758,0
3,11.42,20.38,77.58,386.1,0.1425,0.2839,0.2414,0.1052,0.2597,0.09744,...,26.5,98.87,567.7,0.2098,0.8663,0.6869,0.2575,0.6638,0.173,0
4,20.29,14.34,135.1,1297.0,0.1003,0.1328,0.198,0.1043,0.1809,0.05883,...,16.67,152.2,1575.0,0.1374,0.205,0.4,0.1625,0.2364,0.07678,0


Everything looks as we'd expect and the values are all in the right format 🥳.

We're now at the point where can bring this step (and some others) together into a single
pipeline, the top-level organising entity for code in ZenML. Creating such a pipeline is
as simple as adding a `@pipeline` decorator to a function. This specific
pipeline doesn't return a value, but that option is available to you if you need.

In [5]:
@pipeline
def feature_engineering(
    test_size: float = 0.3,
    drop_na: Optional[bool] = None,
    normalize: Optional[bool] = None,
    drop_columns: Optional[List[str]] = None,
    target: Optional[str] = "target",
    random_state: int = 17
):
    """Feature engineering pipeline."""
    # Link all the steps together by calling them and passing the output
    # of one step as the input of the next step.
    raw_data = data_loader(random_state=random_state, target=target)
    dataset_trn, dataset_tst = data_splitter(
        dataset=raw_data,
        test_size=test_size,
    )
    dataset_trn, dataset_tst, _ = data_preprocessor(
        dataset_trn=dataset_trn,
        dataset_tst=dataset_tst,
        drop_na=drop_na,
        normalize=normalize,
        drop_columns=drop_columns,
        target=target,
        random_state=random_state,
    )

We're ready to run the pipeline now, which we can do just as with the step - by calling the
pipeline function itself:

In [6]:
feature_engineering()

[1;35mInitiating a new run for the pipeline: [0m[1;36mfeature_engineering[1;35m.[0m
[1;35mRegistered new version: [0m[1;36m(version 1)[1;35m.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mStep [0m[1;36mdata_loader[1;35m has started.[0m
[1;35mDataset with 541 records loaded![0m
[1;35mStep [0m[1;36mdata_loader[1;35m has finished in [0m[1;36m0.487s[1;35m.[0m
[1;35mStep [0m[1;36mdata_splitter[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_splitter[1;35m has finished in [0m[1;36m0.920s[1;35m.[0m
[1;35mStep [0m[1;36mdata_preprocessor[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_preprocessor[1;35m has finished in [0m[1;36m0.973s[1;35m.[0m
[1;35mRun [0m[1;36mfeature_engineering-2023_12_14-19_12_49_255201[1;35m has finished in [0m[1;36m2.

Let's run this again with a slightly different test size, to create more datasets:

In [7]:
feature_engineering(test_size=0.25)

[1;35mInitiating a new run for the pipeline: [0m[1;36mfeature_engineering[1;35m.[0m
[1;35mRegistered new version: [0m[1;36m(version 2)[1;35m.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mUsing cached version of [0m[1;36mdata_loader[1;35m.[0m
[1;35mStep [0m[1;36mdata_loader[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_splitter[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_splitter[1;35m has finished in [0m[1;36m0.723s[1;35m.[0m
[1;35mStep [0m[1;36mdata_preprocessor[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_preprocessor[1;35m has finished in [0m[1;36m1.139s[1;35m.[0m
[1;35mRun [0m[1;36mfeature_engineering-2023_12_14-19_12_53_441184[1;35m has finished in [0m[1;36m2.115s[1;35m.[0m
[1;35mYou can visualize your pipeline runs in th

Notice the second time around, the data loader step was **cached**, while the rest of the pipeline was rerun. 
This is because ZenML automatically determined that nothing had changed in the data loader step, 
so it didn't need to rerun it.

Let's run this again with a slightly different test size and random state, to disable the cache and to create more datasets:

In [8]:
feature_engineering(test_size=0.25, random_state=104)

[1;35mInitiating a new run for the pipeline: [0m[1;36mfeature_engineering[1;35m.[0m
[1;35mRegistered new version: [0m[1;36m(version 3)[1;35m.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mStep [0m[1;36mdata_loader[1;35m has started.[0m
[1;35mDataset with 541 records loaded![0m
[1;35mStep [0m[1;36mdata_loader[1;35m has finished in [0m[1;36m0.546s[1;35m.[0m
[1;35mStep [0m[1;36mdata_splitter[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_splitter[1;35m has finished in [0m[1;36m0.717s[1;35m.[0m
[1;35mStep [0m[1;36mdata_preprocessor[1;35m has started.[0m
[1;35mStep [0m[1;36mdata_preprocessor[1;35m has finished in [0m[1;36m0.979s[1;35m.[0m
[1;35mRun [0m[1;36mfeature_engineering-2023_12_14-19_12_56_904810[1;35m has finished in [0m[1;36m2.

At this point you might be interested to view your pipeline runs in the ZenML
Dashboard. In case you are not using a hosted instance of ZenML, you can spin this up by executing the next cell. This will start a
server which you can access by clicking on the link that appears in the output
of the cell.

Log into the Dashboard using default credentials (username 'default' and
password left blank). From there you can inspect the pipeline or the specific
pipeline run.


In [9]:
from zenml.environment import Environment
from zenml.zen_stores.rest_zen_store import RestZenStore


if not isinstance(client.zen_store, RestZenStore):
    # Only spin up a local Dashboard in case you aren't already connected to a remote server
    if Environment.in_google_colab():
        # run ZenML through a cloudflare tunnel to get a public endpoint
        !zenml up --port 8237 & cloudflared tunnel --url http://localhost:8237
    else:
        !zenml up

[1;35mDeploying a local ZenML server with name 'local'.[0m
[?25l[32m⠋[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_server, flavor: local)'.
[2K[1A[2K[1A[2K[32m⠙[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_server, flavor: local)'.
[2K[1A[2K[1A[2K[32m⠹[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_server, flavor: local)'.
[2K[1A[2K[1A[2K[32m⠸[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_server, flavor: local)'.
[2K[1A[2K[1A[2K[32m⠼[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_server, flavor: local)'.
[2K[1A[2K[1A[2K[32m⠴[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_server, flavor: local)'.
[2K[1A[2K[1A[2K[32m⠦[0m Starting service 'LocalZenServer[67a95a27-d110-4e87-a6f1-a1d374144cb9] (type: 
zen_serve

We can also fetch the pipeline from the server and view the results directly in the notebook:

In [10]:
client = Client()
run = client.get_pipeline("feature_engineering").last_run
print(run.name)

feature_engineering-2023_12_14-19_12_56_904810


We can also see the data artifacts that were produced by the last step of the pipeline:

In [11]:
run.steps["data_preprocessor"].outputs

{'dataset_trn': ArtifactVersionResponse(id=UUID('403585cd-2076-413e-a733-9694fcb571bf'), permission_denied=False, body=ArtifactVersionResponseBody(created=datetime.datetime(2023, 12, 14, 19, 12, 58, 668117), updated=datetime.datetime(2023, 12, 14, 19, 12, 58, 668122), user=UserResponse(id=UUID('b11443aa-7044-4002-a761-fca964afead0'), permission_denied=False, body=UserResponseBody(created=datetime.datetime(2023, 12, 14, 19, 12, 36, 615395), updated=datetime.datetime(2023, 12, 14, 19, 12, 36, 615400), active=True, activation_token=None, full_name='', email_opted_in=None, is_service_account=False), metadata=None, name='default'), artifact=ArtifactResponse(id=UUID('6aabc759-379b-488e-aaa2-2e12970aa9eb'), permission_denied=False, body=ArtifactResponseBody(created=datetime.datetime(2023, 12, 14, 19, 12, 51, 32970), updated=datetime.datetime(2023, 12, 14, 19, 12, 51, 32975)), metadata=None, name='dataset_trn'), version='3', uri='/home/apenner/.config/zenml/local_stores/0dd2527d-9162-4251-b3b7

In [12]:
# Read one of the datasets. This is the one with a 0.25 test split
run.steps["data_preprocessor"].outputs["dataset_trn"].load()

Unnamed: 0,mean radius,mean texture,mean perimeter,mean area,mean smoothness,mean compactness,mean concavity,mean concave points,mean symmetry,mean fractal dimension,...,worst texture,worst perimeter,worst area,worst smoothness,worst compactness,worst concavity,worst concave points,worst symmetry,worst fractal dimension,target
444,12.040,28.14,76.85,449.9,0.08752,0.06000,0.023670,0.02377,0.1854,0.05698,...,33.33,87.24,567.6,0.10410,0.09726,0.05524,0.05547,0.2404,0.06639,1
478,9.676,13.14,64.12,272.5,0.12550,0.22040,0.118800,0.07038,0.2057,0.09575,...,18.04,69.47,328.1,0.20060,0.36630,0.29130,0.10750,0.2848,0.13640,1
210,10.440,15.46,66.62,329.6,0.10530,0.07722,0.006643,0.01216,0.1788,0.06450,...,19.80,73.47,395.4,0.13410,0.11530,0.02639,0.04464,0.2615,0.08269,1
299,12.430,17.00,78.60,477.3,0.07557,0.03454,0.013420,0.01699,0.1472,0.05561,...,20.21,81.76,515.9,0.08409,0.04712,0.02237,0.02832,0.1901,0.05932,1
513,14.470,24.99,95.81,656.4,0.08837,0.12300,0.100900,0.03890,0.1872,0.06341,...,31.73,113.50,808.9,0.13400,0.42020,0.40400,0.12050,0.3187,0.10230,1
...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
71,16.070,19.65,104.10,817.7,0.09168,0.08424,0.097690,0.06638,0.1798,0.05391,...,24.56,128.80,1223.0,0.15000,0.20450,0.28290,0.15200,0.2650,0.06387,0
106,10.510,20.19,68.64,334.2,0.11220,0.13030,0.064760,0.03068,0.1922,0.07782,...,22.75,72.62,374.4,0.13000,0.20490,0.12950,0.06136,0.2383,0.09026,1
270,14.410,19.73,96.03,651.0,0.08757,0.16760,0.136200,0.06602,0.1714,0.07192,...,22.13,101.70,767.3,0.09983,0.24720,0.22200,0.10210,0.2272,0.08799,1
435,27.420,26.27,186.90,2501.0,0.10840,0.19880,0.363500,0.16890,0.2061,0.05623,...,31.37,251.20,4254.0,0.13570,0.42560,0.68330,0.26250,0.2641,0.07427,0


We can also get the artifacts directly. Each time you create a new pipeline run, a new `artifact version` is created.

You can fetch these artifact and their versions using the `client`: 

In [13]:
# Get artifact version from our run
dataset_trn_artifact_version_via_run = run.steps["data_preprocessor"].outputs["dataset_trn"] 

# Get latest version from client directly
dataset_trn_artifact_version = client.get_artifact_version("dataset_trn")

# This should be true if our run is the latest run and no artifact has been produced
#  in the intervening time
dataset_trn_artifact_version_via_run.id == dataset_trn_artifact_version.id

True

In [14]:
# Fetch the rest of the artifacts
dataset_tst_artifact_version = client.get_artifact_version("dataset_tst")
preprocessing_pipeline_artifact_version = client.get_artifact_version("preprocess_pipeline")

If you started with a fresh install, then you would have two versions corresponding
to the two pipelines that we ran above. We can even load a artifact version in memory:   

In [15]:
# Load an artifact to verify you can fetch it
dataset_trn_artifact_version.load()

Unnamed: 0,mean radius,mean texture,mean perimeter,mean area,mean smoothness,mean compactness,mean concavity,mean concave points,mean symmetry,mean fractal dimension,...,worst texture,worst perimeter,worst area,worst smoothness,worst compactness,worst concavity,worst concave points,worst symmetry,worst fractal dimension,target
444,12.040,28.14,76.85,449.9,0.08752,0.06000,0.023670,0.02377,0.1854,0.05698,...,33.33,87.24,567.6,0.10410,0.09726,0.05524,0.05547,0.2404,0.06639,1
478,9.676,13.14,64.12,272.5,0.12550,0.22040,0.118800,0.07038,0.2057,0.09575,...,18.04,69.47,328.1,0.20060,0.36630,0.29130,0.10750,0.2848,0.13640,1
210,10.440,15.46,66.62,329.6,0.10530,0.07722,0.006643,0.01216,0.1788,0.06450,...,19.80,73.47,395.4,0.13410,0.11530,0.02639,0.04464,0.2615,0.08269,1
299,12.430,17.00,78.60,477.3,0.07557,0.03454,0.013420,0.01699,0.1472,0.05561,...,20.21,81.76,515.9,0.08409,0.04712,0.02237,0.02832,0.1901,0.05932,1
513,14.470,24.99,95.81,656.4,0.08837,0.12300,0.100900,0.03890,0.1872,0.06341,...,31.73,113.50,808.9,0.13400,0.42020,0.40400,0.12050,0.3187,0.10230,1
...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
71,16.070,19.65,104.10,817.7,0.09168,0.08424,0.097690,0.06638,0.1798,0.05391,...,24.56,128.80,1223.0,0.15000,0.20450,0.28290,0.15200,0.2650,0.06387,0
106,10.510,20.19,68.64,334.2,0.11220,0.13030,0.064760,0.03068,0.1922,0.07782,...,22.75,72.62,374.4,0.13000,0.20490,0.12950,0.06136,0.2383,0.09026,1
270,14.410,19.73,96.03,651.0,0.08757,0.16760,0.136200,0.06602,0.1714,0.07192,...,22.13,101.70,767.3,0.09983,0.24720,0.22200,0.10210,0.2272,0.08799,1
435,27.420,26.27,186.90,2501.0,0.10840,0.19880,0.363500,0.16890,0.2061,0.05623,...,31.37,251.20,4254.0,0.13570,0.42560,0.68330,0.26250,0.2641,0.07427,0


We'll use these artifacts from above in our next pipeline

# ⌚ Step 2: Training pipeline

Now that we have our data it makes sense to train some models to get a sense of
how difficult the task is. The Breast Cancer dataset is sufficiently large and complex 
that it's unlikely we'll be able to train a model that behaves perfectly since the problem 
is inherently complex, but we can get a sense of what a reasonable baseline looks like.

We'll start with two simple models, a SGD Classifier and a Random Forest
Classifier, both batteries-included from `sklearn`. We'll train them both on the
same data and then compare their performance.

<img src=".assets/training_pipeline.png" width="30%" alt="Training pipeline">

In [16]:
import pandas as pd
from sklearn.base import ClassifierMixin
from sklearn.ensemble import RandomForestClassifier
from sklearn.linear_model import SGDClassifier
from typing_extensions import Annotated
from zenml import ArtifactConfig, step
from zenml.logger import get_logger

logger = get_logger(__name__)


@step
def model_trainer(
    dataset_trn: pd.DataFrame,
    model_type: str = "sgd",
) -> Annotated[ClassifierMixin, ArtifactConfig(name="model", is_model_artifact=True)]:
    """Configure and train a model on the training dataset."""
    target = "target"
    if model_type == "sgd":
        model = SGDClassifier()
    elif model_type == "rf":
        model = RandomForestClassifier()
    else:
        raise ValueError(f"Unknown model type {model_type}")   

    logger.info(f"Training model {model}...")

    model.fit(
        dataset_trn.drop(columns=[target]),
        dataset_trn[target],
    )
    return model


Our two training steps both return different kinds of `sklearn` classifier
models, so we use the generic `ClassifierMixin` type hint for the return type.

ZenML allows you to load any version of any dataset that is tracked by the framework
directly into a pipeline using the `ExternalArtifact` interface. This is very convenient
in this case, as we'd like to send our preprocessed dataset from the older pipeline directly
into the training pipeline.

In [17]:
@pipeline
def training(
    train_dataset_id: Optional[UUID] = None,
    test_dataset_id: Optional[UUID] = None,
    model_type: str = "sgd",
    min_train_accuracy: float = 0.0,
    min_test_accuracy: float = 0.0,
):
    """Model training pipeline.""" 
    if train_dataset_id is None or test_dataset_id is None:
        # If we dont pass the IDs, this will run the feature engineering pipeline   
        dataset_trn, dataset_tst = feature_engineering()
    else:
        # Load the datasets from an older pipeline
        dataset_trn = ExternalArtifact(id=train_dataset_id)
        dataset_tst = ExternalArtifact(id=test_dataset_id) 

    trained_model = model_trainer(
        dataset_trn=dataset_trn,
        model_type=model_type,
    )

    model_evaluator(
        model=trained_model,
        dataset_trn=dataset_trn,
        dataset_tst=dataset_tst,
        min_train_accuracy=min_train_accuracy,
        min_test_accuracy=min_test_accuracy,
    )

The end goal of this quick baseline evaluation is to understand which of the two
models performs better. We'll use the `evaluator` step to compare the two
models. This step takes in the model from the trainer step, and computes its score
over the testing set.

In [18]:
# Use a random forest model with the chosen datasets.
# We need to pass the ID's of the datasets into the function
training(
    model_type="rf",
    train_dataset_id=dataset_trn_artifact_version.id,
    test_dataset_id=dataset_tst_artifact_version.id
)

rf_run = client.get_pipeline("training").last_run

[1;35mInitiating a new run for the pipeline: [0m[1;36mtraining[1;35m.[0m
[1;35mRegistered new version: [0m[1;36m(version 1)[1;35m.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mStep [0m[1;36mmodel_trainer[1;35m has started.[0m
[1;35mTraining model RandomForestClassifier()...[0m
[1;35mTraining model RandomForestClassifier()...[0m
  if not hasattr(array, "sparse") and array.dtypes.apply(is_sparse).any():
  if is_sparse(pd_dtype):
  if is_sparse(pd_dtype) or not is_extension_array_dtype(pd_dtype):
  if is_sparse(pd_dtype):
  if is_sparse(pd_dtype) or not is_extension_array_dtype(pd_dtype):
[1;35mStep [0m[1;36mmodel_trainer[1;35m has finished in [0m[1;36m0.487s[1;35m.[0m
[1;35mStep [0m[1;36mmodel_evaluator[1;35m has started.[0m
  if not hasattr(array, "sparse"

In [19]:
# Use a SGD classifier
sgd_run = training(
    model_type="sgd",
    train_dataset_id=dataset_trn_artifact_version.id,
    test_dataset_id=dataset_tst_artifact_version.id
)

sgd_run = client.get_pipeline("training").last_run

[1;35mInitiating a new run for the pipeline: [0m[1;36mtraining[1;35m.[0m
[1;35mRegistered new version: [0m[1;36m(version 2)[1;35m.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mStep [0m[1;36mmodel_trainer[1;35m has started.[0m
[1;35mTraining model SGDClassifier()...[0m
[1;35mTraining model SGDClassifier()...[0m
  if is_sparse(pd_dtype):
  if is_sparse(pd_dtype) or not is_extension_array_dtype(pd_dtype):
  if not hasattr(array, "sparse") and array.dtypes.apply(is_sparse).any():
  if is_sparse(pd_dtype):
  if is_sparse(pd_dtype) or not is_extension_array_dtype(pd_dtype):
[1;35mStep [0m[1;36mmodel_trainer[1;35m has finished in [0m[1;36m0.314s[1;35m.[0m
[1;35mStep [0m[1;36mmodel_evaluator[1;35m has started.[0m
  if not hasattr(array, "sparse") and array.dtypes

You can see from the logs already how our model training went: the
`RandomForestClassifier` performed considerably better than the `SGDClassifier`.
We can use the ZenML `Client` to verify this:

In [20]:
# The evaluator returns a float value with the accuracy
rf_run.steps["model_evaluator"].output.load() > sgd_run.steps["model_evaluator"].output.load()

True

# 💯 Step 3: Associating a model with your pipeline

You can see it is relatively easy to train ML models using ZenML pipelines. But it can be somewhat clunky to track
all the models produced as you develop your experiments and use-cases. Luckily, ZenML offers a *Model Control Plane*,
which is a central register of all your ML models.

You can easily create a ZenML `Model` and associate it with your pipelines using the `ModelVersion` object:

In [21]:
pipeline_settings = {}

# Lets add some metadata to the model to make it identifiable
pipeline_settings["model_version"] = ModelVersion(
    name="breast_cancer_classifier",
    license="Apache 2.0",
    description="A breast cancer classifier",
    tags=["breast_cancer", "classifier"],
)

In [22]:
# Let's train the SGD model and set the version name to "sgd"
pipeline_settings["model_version"].version = "sgd"

# the `with_options` method allows us to pass in pipeline settings
#  and returns a configured pipeline
training_configured = training.with_options(**pipeline_settings)

# We can now run this as usual
training_configured(
    model_type="sgd",
    train_dataset_id=dataset_trn_artifact_version.id,
    test_dataset_id=dataset_tst_artifact_version.id
)

[1;35mInitiating a new run for the pipeline: [0m[1;36mtraining[1;35m.[0m
[1;35mReusing registered version: [0m[1;36m(version: 2)[1;35m.[0m
[1;35mNew model [0m[1;36mbreast_cancer_classifier[1;35m was created implicitly.[0m
[1;35mNew model version [0m[1;36msgd[1;35m was created.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mUsing cached version of [0m[1;36mmodel_trainer[1;35m.[0m
[1;35mStep [0m[1;36mmodel_trainer[1;35m has started.[0m
[1;35mUsing cached version of [0m[1;36mmodel_evaluator[1;35m.[0m
[1;35mLinking artifact [0m[1;36moutput[1;35m to model [0m[1;36mNone[1;35m version [0m[1;36mNone[1;35m implicitly.[0m
[1;35mStep [0m[1;36mmodel_evaluator[1;35m has started.[0m
[1;35mRun [0m[1;36mtraining-2023_12_14-19_13_28_867038[1;35m has f

In [23]:
# Let's train the RF model and set the version name to "rf"
pipeline_settings["model_version"].version = "rf"

# the `with_options` method allows us to pass in pipeline settings
#  and returns a configured pipeline
training_configured = training.with_options(**pipeline_settings)

# Let's run it again to make sure we have two versions
training_configured(
    model_type="rf",
    train_dataset_id=dataset_trn_artifact_version.id,
    test_dataset_id=dataset_tst_artifact_version.id
)

[1;35mInitiating a new run for the pipeline: [0m[1;36mtraining[1;35m.[0m
[1;35mReusing registered version: [0m[1;36m(version: 1)[1;35m.[0m
[1;35mNew model version [0m[1;36mrf[1;35m was created.[0m
[1;35mExecuting a new run.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mUsing cached version of [0m[1;36mmodel_trainer[1;35m.[0m
[1;35mStep [0m[1;36mmodel_trainer[1;35m has started.[0m
[1;35mUsing cached version of [0m[1;36mmodel_evaluator[1;35m.[0m
[1;35mLinking artifact [0m[1;36moutput[1;35m to model [0m[1;36mNone[1;35m version [0m[1;36mNone[1;35m implicitly.[0m
[1;35mStep [0m[1;36mmodel_evaluator[1;35m has started.[0m
[1;35mRun [0m[1;36mtraining-2023_12_14-19_13_30_626808[1;35m has finished in [0m[1;36m0.321s[1;35m.[0m
[1;35mYou can visualize your pipeline runs in t

This time, running both pipelines has created two associated **model versions**.
You can list your ZenML model and their versions as follows:

In [24]:
zenml_model = client.get_model("breast_cancer_classifier")
print(zenml_model)

print(f"Model {zenml_model.name} has {len(zenml_model.versions)} versions")

zenml_model.versions[0].version, zenml_model.versions[1].version

id=UUID('aa651410-ae18-437b-b55b-81ed7d85e093') permission_denied=False body=ModelResponseBody(created=datetime.datetime(2023, 12, 14, 19, 13, 28, 743771), updated=datetime.datetime(2023, 12, 14, 19, 13, 28, 743774), user=UserResponse(id=UUID('b11443aa-7044-4002-a761-fca964afead0'), permission_denied=False, body=UserResponseBody(created=datetime.datetime(2023, 12, 14, 19, 12, 36, 615395), updated=datetime.datetime(2023, 12, 14, 19, 12, 36, 615400), active=True, activation_token=None, full_name='', email_opted_in=None, is_service_account=False), metadata=None, name='default'), tags=[TagResponseModel(id=UUID('f84d2acb-50c0-45aa-9e31-b0b2b5cadeb1'), created=datetime.datetime(2023, 12, 14, 19, 13, 28, 747775), updated=datetime.datetime(2023, 12, 14, 19, 13, 28, 747778), missing_permissions=False, name='breast_cancer', color=<ColorVariants.MAGENTA: 'magenta'>, tagged_count=1), TagResponseModel(id=UUID('217927ee-4157-441a-93b3-151ad3724ebb'), created=datetime.datetime(2023, 12, 14, 19, 13, 2

('sgd', 'rf')

The interesting part is that ZenML went ahead and linked all artifacts produced by the
pipelines to that model version, including the two pickle files that represent our
SGD and RandomForest classifier. We can see all artifacts directly from the model
version object:

In [25]:
# Let's load the RF version
rf_zenml_model_version = client.get_model_version("breast_cancer_classifier", "rf")

# We can now load our classifier directly as well
random_forest_classifier = rf_zenml_model_version.get_artifact("model").load()

random_forest_classifier

If you are a [ZenML Cloud](https://zenml.io/cloud) user, you can see all of this visualized in the dashboard:

<img src=".assets/cloud_mcp_screenshot.png" width="70%" alt="Model Control Plane">

There is a lot more you can do with ZenML models, including the ability to
track metrics by adding metadata to it, or having them persist in a model
registry. However, these topics can be explored more in the
[ZenML docs](https://docs.zenml.io).

For now, we will use the ZenML model control plane to promote our best
model to `production`. You can do this by simply setting the `stage` of
your chosen model version to the `production` tag.

In [26]:
# Set our best classifier to production
rf_zenml_model_version.set_stage("production", force=True)

Of course, normally one would only promote the model by comparing to all other model
versions and doing some other tests. But that's a bit more advanced use-case. See the
[e2e_batch example](https://github.com/zenml-io/zenml/tree/main/examples/e2e) to get
more insight into that sort of flow!

<img src=".assets/cloud_mcp.png" width="60%" alt="Model Control Plane">

Once the model is promoted, we can now consume the right model version in our
batch inference pipeline directly. Let's see how that works.

# 🫅 Step 4: Consuming the model in production

The batch inference pipeline simply takes the model marked as `production` and runs inference on it
with `live data`. The critical step here is the `inference_predict` step, where we load the model in memory
and generate predictions:

<img src=".assets/inference_pipeline.png" width="45%" alt="Inference pipeline">

In [27]:
@step
def inference_predict(dataset_inf: pd.DataFrame) -> Annotated[pd.Series, "predictions"]:
    """Predictions step"""
    # Get the model_version
    model_version = get_step_context().model_version

    # run prediction from memory
    predictor = model_version.load_artifact("model")
    predictions = predictor.predict(dataset_inf)

    predictions = pd.Series(predictions, name="predicted")

    return predictions


Apart from the loading the model, we must also load the preprocessing pipeline that we ran in feature engineering,
so that we can do the exact steps that we did on training time, in inference time. Let's bring it all together:

In [28]:
@pipeline
def inference(preprocess_pipeline_id: UUID):
    """Model batch inference pipeline"""
    # random_state = client.get_artifact_version(id=preprocess_pipeline_id).metadata["random_state"].value
    # target = client.get_artifact_version(id=preprocess_pipeline_id).run_metadata['target'].value
    random_state = 42
    target = "target"

    df_inference = data_loader(
        random_state=random_state, is_inference=True
    )
    df_inference = inference_preprocessor(
        dataset_inf=df_inference,
        # We use the preprocess pipeline from the feature engineering pipeline
        preprocess_pipeline=ExternalArtifact(id=preprocess_pipeline_id),
        target=target,
    )
    inference_predict(
        dataset_inf=df_inference,
    )


The way to load the right model is to pass in the `production` stage into the `ModelVersion` config this time.
This will ensure to always load the production model, decoupled from all other pipelines:

In [29]:
pipeline_settings = {"enable_cache": False}

# Lets add some metadata to the model to make it identifiable
pipeline_settings["model_version"] = ModelVersion(
    name="breast_cancer_classifier",
    version="production", # We can pass in the stage name here!
    license="Apache 2.0",
    description="A breast cancer classifier",
    tags=["breast_cancer", "classifier"],
)

[1;35m[0m[1;36mversion[1;35m [0m[1;36mproduction[1;35m matches one of the possible [0m[1;36mModelStages[1;35m and will be fetched using stage.[0m


In [30]:
# the `with_options` method allows us to pass in pipeline settings
#  and returns a configured pipeline
inference_configured = inference.with_options(**pipeline_settings)

# Let's run it again to make sure we have two versions
# We need to pass in the ID of the preprocessing done in the feature engineering pipeline
# in order to avoid training-serving skew
inference_configured(
    preprocess_pipeline_id=preprocessing_pipeline_artifact_version.id
)

[1;35mInitiating a new run for the pipeline: [0m[1;36minference[1;35m.[0m
[1;35mRegistered new version: [0m[1;36m(version 1)[1;35m.[0m
[1;35mExecuting a new run.[0m
[1;35mCaching is disabled by default for [0m[1;36minference[1;35m.[0m
[1;35mUsing user: [0m[1;36mdefault[1;35m[0m
[1;35mUsing stack: [0m[1;36mdefault[1;35m[0m
[1;35m  artifact_store: [0m[1;36mdefault[1;35m[0m
[1;35m  orchestrator: [0m[1;36mdefault[1;35m[0m
[1;35mStep [0m[1;36mdata_loader[1;35m has started.[0m
[1;35mDataset with 28 records loaded![0m
[1;35mStep [0m[1;36mdata_loader[1;35m has finished in [0m[1;36m0.642s[1;35m.[0m
[1;35mStep [0m[1;36minference_preprocessor[1;35m has started.[0m
[1;35mStep [0m[1;36minference_preprocessor[1;35m has finished in [0m[1;36m0.847s[1;35m.[0m
[1;35mStep [0m[1;36minference_predict[1;35m has started.[0m
[33mYou specified both an ID as well as a version of the artifact_versions. Ignoring the version and fetching the ar

ZenML automatically links all artifacts to the `production` model version as well, including the predictions
that were returned in the pipeline. This completes the MLOps loop of training to inference:

In [31]:
# Fetch production model
production_model_version = client.get_model_version("breast_cancer_classifier", "production")

# Get the predictions artifact
production_model_version.get_artifact("predictions").load()

0     1
1     0
2     0
3     1
4     1
5     0
6     0
7     0
8     0
9     1
10    1
11    0
12    1
13    0
14    1
15    0
16    1
17    1
18    1
19    0
20    1
21    1
22    0
23    1
24    1
25    1
26    1
27    1
Name: series, dtype: int64

You can also see all predictions ever created as a complete history in the dashboard:

<img src=".assets/cloud_mcp_predictions.png" width="70%" alt="Model Control Plane">

## Congratulations!

You're a legit MLOps engineer now! You trained two models, evaluated them against
a test set, registered the best one with the ZenML model control plane,
and served some predictions. You also learned how to iterate on your models and
data by using some of the ZenML utility abstractions. You saw how to view your
artifacts and stacks via the client as well as the ZenML Dashboard.

## Further exploration

This was just the tip of the iceberg of what ZenML can do; check out the [**docs**](https://docs.zenml.io/) to learn more
about the capabilities of ZenML. For example, you might want to:

- Run the same pipeline on a cloud stack in production.
- Track your metrics in an experiment tracker like [MLflow]().
- Learn how to transition your code from this notebook setting to a production setting.

## What next?

* If you have questions or feedback... join our [**Slack Community**](https://zenml.io/slack) and become part of the ZenML family!
* If you want to quickly get started with ZenML, check out the [ZenML Cloud](https://zenml.io/cloud).