ZenML: Modular and Cloud-Agnostic MLOps Pipeline Framework
Introduction
Building an accurate machine learning model is only a small part of the challenge in ML production. The real challenge lies in creating reproducible, scalable, and maintainable pipelines. ZenML is an open-source MLOps framework designed to tackle these problems with a modular and cloud-agnostic approach.
ZenML allows data scientists and ML engineers to define ML pipelines using simple Python decorators while providing full flexibility to integrate various tools and cloud platforms. With its unique "stacks" concept, you can switch from local development to cloud production without changing your pipeline code.
In this tutorial, we will learn ZenML from the basics to building an end-to-end pipeline covering data ingestion, preprocessing, training, evaluation, and deployment.
Prerequisites
Before getting started, make sure you have:
- Python 3.8 or later
- pip package manager
- Basic understanding of machine learning
- Familiarity with scikit-learn or other ML frameworks
Installing ZenML
Installing ZenML is straightforward using pip:
pip install zenml
For installation with additional integrations:
pip install "zenml[server]"
After installation, initialize a ZenML repository in your project:
zenml init
Launch the ZenML dashboard locally:
zenml login --local
The dashboard will be available at http://127.0.0.1:8237, providing complete visualization of your pipelines, artifacts, and stacks.
Core Concepts: Steps and Pipelines
Defining Steps with @step
A step is the smallest unit of work in ZenML. Each step is a Python function decorated with @step:
from zenml import step
import pandas as pd
@step
def loaddata() -> pd.DataFrame:
"""Load dataset from data source."""
df = pd.readcsv("data/trainingdata.csv")
return df
@step
def preprocessdata(df: pd.DataFrame) -> pd.DataFrame:
"""Clean and process data."""
df = df.dropna()
df = df.dropduplicates()
return df
ZenML automatically performs type checking and serialization on each step's input/output based on the type hints you provide.
Defining Pipelines with @pipeline
A pipeline connects multiple steps into a structured workflow:
from zenml import pipeline
@pipeline
def trainingpipeline():
"""Pipeline for ML model training."""
data = loaddata()
processeddata = preprocessdata(data)
model = trainmodel(processeddata)
metrics = evaluatemodel(model, processeddata)
return metrics
Running a pipeline is as simple as calling the function:
if name == "main":
trainingpipeline()
Parameterizing Steps
You can create steps that accept configuration parameters:
from zenml import step
from pydantic import BaseModel
class TrainingConfig(BaseModel):
learningrate: float = 0.01
nestimators: int = 100
maxdepth: int = 5
@step
def trainmodel(
data: pd.DataFrame,
config: TrainingConfig = TrainingConfig()
) -> object:
"""Train model with customizable configuration."""
from sklearn.ensemble import RandomForestClassifier
X = data.drop("target", axis=1)
y = data["target"]
model = RandomForestClassifier(
nestimators=config.nestimators,
maxdepth=config.maxdepth
)
model.fit(X, y)
return model
Artifacts and Materializers
Understanding Artifacts
Every output from a step is automatically stored as an artifact. Artifacts are data produced and consumed by steps in a pipeline. ZenML tracks each artifact with comprehensive metadata including version, type, and lineage.
from zenml import step, logartifactmetadata
@step
def trainmodel(data: pd.DataFrame) -> object:
model = RandomForestClassifier()
model.fit(Xtrain, ytrain)
logartifactmetadata(
artifactname="model",
metadata={
"accuracy": float(model.score(Xtest, ytest)),
"nfeatures": Xtrain.shape[1],
"algorithm": "RandomForest"
}
)
return model
Custom Materializers
Materializers control how artifacts are serialized and deserialized. You can create custom materializers for specialized data types:
from zenml.materializers import BaseMaterializer
from zenml.enums import ArtifactType
import json
import os
class CustomModelMaterializer(BaseMaterializer):
ASSOCIATEDTYPES = (CustomModel,)
ASSOCIATEDARTIFACTTYPE = ArtifactType.MODEL
def load(self, datatype):
filepath = os.path.join(self.uri, "model.json")
with open(filepath, "r") as f:
data = json.load(f)
return CustomModel.fromdict(data)
def save(self, model):
filepath = os.path.join(self.uri, "model.json")
with open(filepath, "w") as f:
json.dump(model.todict(), f)
Stacks: Modular Infrastructure
The Stack Concept
A stack is an infrastructure configuration that defines where and how pipelines are executed. Each stack consists of several components:
- Orchestrator: Runs pipelines (local, Airflow, Kubeflow, etc.)
- Artifact Store: Stores artifacts (local, S3, GCS, Azure Blob)
- Experiment Tracker: Tracks experiments (MLflow, WandB, Neptune)
- Model Deployer: Deploys models (MLflow, Seldon, BentoML)
- Container Registry: Stores Docker images
- Step Operator: Runs steps in specialized environments (SageMaker, Vertex AI)
Managing Stacks
# List available stacks
zenml stack list
Register new components
zenml artifact-store register mys3store \
--flavor=s3 \
--path=s3://my-bucket/zenml
zenml orchestrator register mykubeflow \
--flavor=kubeflow \
--kubernetescontext=my-cluster
Create a new stack
zenml stack register productionstack \
--orchestrator=mykubeflow \
--artifact-store=mys3store
Set active stack
zenml stack set productionstack
Stack Components in Detail
Here is a more complete stack configuration example:
# Register MLflow experiment tracker
zenml experiment-tracker register mlflowtracker \
--flavor=mlflow \
--trackinguri=http://mlflow-server:5000
Register model deployer
zenml model-deployer register mlflowdeployer \
--flavor=mlflow
Complete production stack
zenml stack register fullstack \
--orchestrator=mykubeflow \
--artifact-store=mys3store \
--experiment-tracker=mlflowtracker \
--model-deployer=mlflowdeployer
Integration with MLflow and Weights & Biases
MLflow Integration
Install the MLflow integration:
zenml integration install mlflow -y
Use the MLflow experiment tracker in a step:
from zenml import step
from zenml.integrations.mlflow.flavors.mlflowexperimenttrackerflavor import (
MLFlowExperimentTrackerSettings
)
import mlflow
mlflowsettings = MLFlowExperimentTrackerSettings(
experimentname="myexperiment"
)
@step(
experimenttracker="mlflowtracker",
settings={"experimenttracker": mlflowsettings}
)
def trainwithmlflow(data: pd.DataFrame) -> object:
mlflow.autolog()
model = RandomForestClassifier(nestimators=100)
model.fit(Xtrain, ytrain)
accuracy = model.score(Xtest, ytest)
mlflow.logmetric("accuracy", accuracy)
mlflow.logparam("nestimators", 100)
return model
Weights & Biases Integration
zenml integration install wandb -y
from zenml.integrations.wandb.flavors.wandbexperimenttrackerflavor import (
WandbExperimentTrackerSettings
)
import wandb
wandb
settings = WandbExperimentTrackerSettings(
project="my-zenml-project",
entity="my-team"
)
@step(
experimenttracker="wandbtracker",
settings={"experimenttracker": wandbsettings}
)
def trainwithwandb(data: pd.DataFrame) -> object:
config = {"nestimators": 100, "maxdepth": 5}
wandb.config.update(config)
model = RandomForestClassifier(*config)
model.fit(Xtrain, ytrain)
wandb.log({"accuracy": model.score(Xtest, ytest)})
return model
Model Deployment
ZenML supports model deployment through various deployers:
from zenml import step, pipeline
from zenml.integrations.mlflow.steps import mlflowmodeldeployerstep
@step
def deploymenttrigger(accuracy: float) -> bool:
"""Determine if the model is ready for deployment."""
return accuracy > 0.85
@pipeline
def deploymentpipeline():
data = loaddata()
processed = preprocessdata(data)
model = trainmodel(processed)
accuracy = evaluatemodel(model, processed)
shoulddeploy = deploymenttrigger(accuracy)
if shoulddeploy:
mlflowmodeldeployerstep(
model=model,
deploydecision=shoulddeploy,
workers=3
)
Caching
ZenML enables caching by default for every step. If a step's input and code haven't changed, ZenML will reuse the previous result:
# Disable caching for a specific step
@step(enablecache=False)
def alwaysfreshdata() -> pd.DataFrame:
"""This step always re-executes."""
return pd.readcsv("livedata.csv")
Disable caching for the entire pipeline
@pipeline(enablecache=False)
def nocachepipeline():
data = alwaysfreshdata()
process(data)
Caching is extremely useful for saving time and resources when iterating on specific parts of the pipeline.
Pipeline Scheduling
ZenML supports periodic pipeline scheduling:
from zenml.config.schedule import Schedule
Schedule pipeline to run daily
schedule = Schedule(
cronexpression="0 8 " # Every day at 8 AM
)
trainingpipeline = trainingpipeline.withoptions(
schedule=schedule
)
trainingpipeline()
You can also schedule with intervals:
from datetime import timedelta
schedule = Schedule(
intervalsecond=timedelta(hours=6),
starttime="2026-01-01T00:00:00"
)
ZenML Dashboard
ZenML provides a powerful web dashboard for monitoring:
# Login to ZenML server
zenml login --local
Or connect to a remote server
zenml login https://my-zenml-server.com
The dashboard provides:
- Pipeline Runs: DAG visualization and status of each step
- Artifacts: Browser for viewing all generated artifacts
- Stacks: Stack and infrastructure component management
- Model Registry: Model version tracking and deployment status
Practical Example: End-to-End ML Pipeline
Here is a complete pipeline example from data ingestion to deployment:
import pandas as pd
import numpy as np
from sklearn.modelselection import traintestsplit
from sklearn.preprocessing import StandardScaler, LabelEncoder
from sklearn.ensemble import RandomForestClassifier
from sklearn.metrics import accuracyscore, classificationreport
from zenml import step, pipeline, logartifactmetadata, ArtifactConfig
from zenml.client import Client
from typing import Tuple
from typingextensions import Annotated
from pydantic import BaseModel
Configuration
class ModelConfig(BaseModel):
testsize: float = 0.2
nestimators: int = 200
maxdepth: int = 10
randomstate: int = 42
minaccuracy: float = 0.85
Step 1: Data Ingestion
@step
def ingestdata() -> Annotated[pd.DataFrame, "rawdata"]:
"""Fetch data from source."""
from sklearn.datasets import loadiris
iris = loadiris(asframe=True)
df = iris.frame
df.columns = ["sepallength", "sepalwidth",
"petallength", "petalwidth", "target"]
logartifactmetadata(
artifactname="rawdata",
metadata={
"numrows": len(df),
"numcolumns": len(df.columns),
"columns": list(df.columns)
}
)
return df
Step 2: Data Preprocessing
@step
def preprocess(
df: pd.DataFrame,
config: ModelConfig = ModelConfig()
) -> Tuple[
Annotated[np.ndarray, "Xtrain"],
Annotated[np.ndarray, "Xtest"],
Annotated[np.ndarray, "ytrain"],
Annotated[np.ndarray, "ytest"],
Annotated[StandardScaler, "scaler"]
]:
"""Preprocessing: split, scale, encode."""
X = df.drop("target", axis=1).values
y = df["target"].values
Xtrain, Xtest, ytrain, ytest = traintestsplit(
X, y,
testsize=config.testsize,
randomstate=config.randomstate,
stratify=y
)
scaler = StandardScaler()
Xtrain = scaler.fittransform(Xtrain)
Xtest = scaler.transform(Xtest)
return Xtrain, Xtest, ytrain, ytest, scaler
Step 3: Model Training
@step
def train(
Xtrain: np.ndarray,
ytrain: np.ndarray,
config: ModelConfig = ModelConfig()
) -> Annotated[RandomForestClassifier, "trainedmodel"]:
"""Train a RandomForest model."""
model = RandomForestClassifier(
nestimators=config.nestimators,
maxdepth=config.maxdepth,
randomstate=config.randomstate,
njobs=-1
)
model.fit(Xtrain, ytrain)
logartifactmetadata(
artifactname="trainedmodel",
metadata={
"algorithm": "RandomForestClassifier",
"nestimators": config.nestimators,
"maxdepth": config.maxdepth
}
)
return model
Step 4: Model Evaluation
@step
def evaluate(
model: RandomForestClassifier,
Xtest: np.ndarray,
ytest: np.ndarray,
config: ModelConfig = ModelConfig()
) -> Annotated[float, "accuracy"]:
"""Evaluate model performance."""
ypred = model.predict(Xtest)
accuracy = accuracyscore(ytest, ypred)
report = classificationreport(ytest, ypred, outputdict=True)
logartifactmetadata(
artifactname="accuracy",
metadata={
"accuracy": accuracy,
"precisionmacro": report["macro avg"]["precision"],
"recallmacro": report["macro avg"]["recall"],
"f1macro": report["macro avg"]["f1-score"]
}
)
print(f"Model Accuracy: {accuracy:.4f}")
print(classificationreport(ytest, ypred))
return accuracy
Step 5: Deployment Decision
@step
def decidedeployment(
accuracy: float,
config: ModelConfig = ModelConfig()
) -> Annotated[bool, "deploydecision"]:
"""Decide whether the model is ready for deployment."""
shoulddeploy = accuracy >= config.minaccuracy
if shoulddeploy:
print(f"Model APPROVED for deployment (accuracy: {accuracy:.4f} >= {config.minaccuracy})")
else:
print(f"Model REJECTED for deployment (accuracy: {accuracy:.4f} < {config.minaccuracy})")
return shoulddeploy
Pipeline Definition
@pipeline
def mlpipeline():
"""End-to-end ML pipeline."""
rawdata = ingestdata()
Xtrain, Xtest, ytrain, ytest, scaler = preprocess(rawdata)
model = train(Xtrain, ytrain)
accuracy = evaluate(model, Xtest, ytest)
deploydecision = decidedeployment(accuracy)
return deploydecision
Run the pipeline
if name == "main":
run = mlpipeline()
print(f"Pipeline completed! Run ID: {run.id}")
Best Practices
from zenml.config import DockerSettings
dockersettings = DockerSettings(
requirements=["scikit-learn==1.3.0", "pandas==2.0.0"],
requiredintegrations=["mlflow"]
)
@pipeline(settings={"docker": dockersettings})
def reproducible_pipeline():
...
Conclusion
ZenML provides an elegant framework for building reproducible and modular MLOps pipelines. With its flexible stacks concept, you can easily switch between local development and cloud production without changing your pipeline code. Extensive integrations with popular tools like MLflow and Weights & Biases make it a solid choice for teams looking to implement MLOps best practices.
Features like automatic caching, artifact tracking, and pipeline scheduling make ZenML a comprehensive solution for managing the machine learning lifecycle from experimentation to production.