Metaflow Tutorial: Netflix's MLOps Framework for Data Science
Metaflow is an open-source framework originally developed at Netflix for building and managing real-life data science projects. It allows data scientists to focus on modeling without worrying about infrastructure, orchestration, and deployment. In this tutorial, we'll learn how to use Metaflow from basics to advanced features.
Why Metaflow?
In data science and machine learning, one of the biggest challenges isn't building models but managing the entire lifecycle from experimentation to production. Metaflow addresses this with a pragmatic approach:
- Human-centric API: Designed to be easy for data scientists, not just ML engineers
- Automatic versioning: Every experiment is automatically tracked and reproducible
- Transparent scalability: Move from laptop to cloud without changing code
- Dependency management: Manages Python environments automatically
- Integration: Works with AWS, Azure, GCP, and Kubernetes
Installation
Basic Installation
pip install metaflow
Installation with AWS Support
pip install metaflow[aws]
Verify Installation
import metaflow
print(metaflow.version)
Initial Configuration
After installation, run the configuration:
metaflow configure show
To configure AWS S3 as datastore:
metaflow configure aws
Core Concepts: Flow and Step
Metaflow uses Flow and Step concepts to organize data science pipelines. A Flow is a DAG (Directed Acyclic Graph) composed of multiple Steps.
Your First Flow
from metaflow import FlowSpec, step
class HelloFlow(FlowSpec):
@step
def start(self):
print("Starting the first flow!")
self.message = "Hello from Metaflow"
self.next(self.end)
@step
def end(self):
print(f"Message: {self.message}")
print("Flow completed!")
if name == 'main':
HelloFlow()
Run the flow:
python helloflow.py run
Understanding Artifacts
Every variable saved as self.x in a step becomes an artifact that is automatically versioned and can be accessed later.
from metaflow import FlowSpec, step
class ArtifactFlow(FlowSpec):
@step
def start(self):
self.data = [1, 2, 3, 4, 5]
self.modelname = "randomforest"
self.next(self.process)
@step
def process(self):
self.result = sum(self.data) 2
print(f"Model: {self.modelname}")
print(f"Result: {self.result}")
self.next(self.end)
@step
def end(self):
print(f"Final result: {self.result}")
if name == 'main':
ArtifactFlow()
Branching and Join
Metaflow supports parallel execution through branching. This is useful for comparing multiple models simultaneously.
Branching Example
from metaflow import FlowSpec, step
class BranchFlow(FlowSpec):
@step
def start(self):
self.rawdata = list(range(100))
self.next(self.trainrf, self.trainxgb)
@step
def trainrf(self):
from sklearn.ensemble import RandomForestClassifier
from sklearn.datasets import makeclassification
X, y = makeclassification(nsamples=1000, nfeatures=20, randomstate=42)
model = RandomForestClassifier(nestimators=100, randomstate=42)
model.fit(X, y)
self.accuracy = model.score(X, y)
self.modeltype = "RandomForest"
print(f"RF Accuracy: {self.accuracy:.4f}")
self.next(self.join)
@step
def trainxgb(self):
from sklearn.ensemble import GradientBoostingClassifier
from sklearn.datasets import makeclassification
X, y = makeclassification(nsamples=1000, nfeatures=20, randomstate=42)
model = GradientBoostingClassifier(nestimators=100, randomstate=42)
model.fit(X, y)
self.accuracy = model.score(X, y)
self.modeltype = "GradientBoosting"
print(f"XGB Accuracy: {self.accuracy:.4f}")
self.next(self.join)
@step
def join(self, inputs):
best = max(inputs, key=lambda x: x.accuracy)
self.bestmodel = best.modeltype
self.bestaccuracy = best.accuracy
print(f"Best model: {self.bestmodel} ({self.bestaccuracy:.4f})")
self.next(self.end)
@step
def end(self):
print(f"Pipeline complete. Selected model: {self.bestmodel}")
if name == 'main':
BranchFlow()
Parameters and Configuration
Metaflow supports parameters that can be set at runtime without modifying code.
from metaflow import FlowSpec, step, Parameter
class ParamFlow(FlowSpec):
learningrate = Parameter(
'learningrate',
help='Learning rate for training',
default=0.01
)
nestimators = Parameter(
'nestimators',
help='Number of estimators',
default=100,
type=int
)
datapath = Parameter(
'datapath',
help='Path to dataset',
default='data/train.csv'
)
@step
def start(self):
print(f"Learning Rate: {self.learningrate}")
print(f"N Estimators: {self.nestimators}")
print(f"Data Path: {self.datapath}")
self.next(self.train)
@step
def train(self):
from sklearn.ensemble import GradientBoostingClassifier
from sklearn.datasets import makeclassification
X, y = makeclassification(nsamples=1000, randomstate=42)
model = GradientBoostingClassifier(
learningrate=self.learningrate,
nestimators=self.nestimators,
randomstate=42
)
model.fit(X, y)
self.score = model.score(X, y)
print(f"Score: {self.score:.4f}")
self.next(self.end)
@step
def end(self):
print(f"Training complete with score: {self.score:.4f}")
if name == 'main':
ParamFlow()
Run with custom parameters:
python paramflow.py run --learningrate 0.1 --nestimators 200
Foreach: Parallel Processing
foreach allows you to run the same step for multiple items in parallel.
from metaflow import FlowSpec, step
class ForeachFlow(FlowSpec):
@step
def start(self):
self.datasets = ['dataseta', 'datasetb', 'datasetc']
self.next(self.process, foreach='datasets')
@step
def process(self):
datasetname = self.input
print(f"Processing: {datasetname}")
import random
random.seed(hash(datasetname) % 232)
self.accuracy = random.uniform(0.7, 0.99)
self.dataset = datasetname
print(f"{datasetname}: accuracy = {self.accuracy:.4f}")
self.next(self.join)
@step
def join(self, inputs):
self.results = {
inp.dataset: inp.accuracy for inp in inputs
}
best = max(inputs, key=lambda x: x.accuracy)
self.bestdataset = best.dataset
self.bestaccuracy = best.accuracy
print("\nAll dataset results:")
for name, acc in self.results.items():
print(f" {name}: {acc:.4f}")
print(f"\nBest: {self.bestdataset} ({self.bestaccuracy:.4f})")
self.next(self.end)
@step
def end(self):
print("Foreach pipeline complete!")
if name == 'main':
ForeachFlow()
Advanced Features
@conda and @pypi Decorators
Manage per-step dependencies using decorators:
from metaflow import FlowSpec, step, condabase, pypi
@conda
base(python='3.10')
class MLFlow(FlowSpec):
@pypi(packages={'scikit-learn': '1.4.0', 'pandas': '2.2.0'})
@step
def start(self):
import sklearn
import pandas as pd
print(f"sklearn: {sklearn.version}")
print(f"pandas: {pd.version}")
self.next(self.end)
@step
def end(self):
print("Done!")
if name == 'main':
MLFlow()
@resources Decorator
Configure compute resources for each step:
from metaflow import FlowSpec, step, resources
class GPUFlow(FlowSpec):
@resources(cpu=4, memory=8000)
@step
def start(self):
self.data = list(range(1000000))
self.next(self.train)
@resources(cpu=8, memory=16000, gpu=1)
@step
def train(self):
print("Training with GPU...")
self.modeltrained = True
self.next(self.end)
@step
def end(self):
print(f"Model trained: {self.modeltrained}")
if name == 'main':
GPUFlow()
Retry and Error Handling
from metaflow import FlowSpec, step, retry, catch
class RobustFlow(FlowSpec):
@retry(times=3)
@step
def start(self):
self.data = self.loaddata()
self.next(self.train)
def loaddata(self):
import random
if random.random() < 0.3:
raise Exception("Connection lost!")
return [1, 2, 3, 4, 5]
@catch(var='trainerror')
@step
def train(self):
self.result = sum(self.data) 2
self.next(self.end)
@step
def end(self):
if hasattr(self, 'trainerror') and self.trainerror:
print(f"Training failed: {self.trainerror}")
else:
print(f"Result: {self.result}")
if name == 'main':
RobustFlow()
Client API: Accessing Experiment Results
One of the most powerful features of Metaflow is the Client API for accessing previous experiment results:
from metaflow import Flow, Run, Step
flow = Flow('BranchFlow')
for run in flow.runs():
print(f"Run ID: {run.id}")
print(f" Status: {run.successful}")
print(f" Time: {run.createdat}")
if run.successful:
endstep = Step(f'BranchFlow/{run.id}/end')
for task in endstep.tasks():
print(f" Best Model: {task['bestmodel'].data}")
print(f" Best Accuracy: {task['bestaccuracy'].data}")
print()
Accessing the Latest Run
from metaflow import Flow
run = Flow('BranchFlow').latestsuccessfulrun
print(f"Latest run: {run.id}")
print(f"Best model: {run.data.bestmodel}")
print(f"Best accuracy: {run.data.bestaccuracy}")
Complete Project: End-to-End ML Pipeline
Here's a complete ML pipeline using Metaflow:
from metaflow import FlowSpec, step, Parameter, retry, catch
import json
class MLPipelineFlow(FlowSpec):
testsize = Parameter('testsize', default=0.2, type=float)
randomstate = Parameter('randomstate', default=42, type=int)
@step
def start(self):
print("=== ML Pipeline Started ===")
self.next(self.loaddata)
@retry(times=2)
@step
def loaddata(self):
from sklearn.datasets import loadbreastcancer
import pandas as pd
data = loadbreastcancer()
self.df = pd.DataFrame(data.data, columns=data.featurenames)
self.target = list(data.target)
self.featurenames = list(data.featurenames)
print(f"Dataset loaded: {self.df.shape}")
print(f"Features: {len(self.featurenames)}")
self.next(self.preprocess)
@step
def preprocess(self):
from sklearn.modelselection import traintestsplit
from sklearn.preprocessing import StandardScaler
import numpy as np
X = self.df.values
y = np.array(self.target)
Xtrain, Xtest, ytrain, ytest = traintestsplit(
X, y, testsize=self.testsize, randomstate=self.randomstate
)
scaler = StandardScaler()
self.Xtrain = scaler.fittransform(Xtrain).tolist()
self.Xtest = scaler.transform(Xtest).tolist()
self.ytrain = ytrain.tolist()
self.ytest = ytest.tolist()
self.scalermean = scaler.mean.tolist()
self.scalerscale = scaler.scale.tolist()
print(f"Train: {len(self.Xtrain)}, Test: {len(self.Xtest)}")
self.next(self.trainlr, self.trainrf, self.trainsvm)
@step
def trainlr(self):
from sklearn.linearmodel import LogisticRegression
from sklearn.metrics import accuracyscore, f1score
import numpy as np
Xtrain = np.array(self.Xtrain)
ytrain = np.array(self.ytrain)
Xtest = np.array(self.Xtest)
ytest = np.array(self.ytest)
model = LogisticRegression(maxiter=1000, randomstate=self.randomstate)
model.fit(Xtrain, ytrain)
ypred = model.predict(Xtest)
self.accuracy = accuracyscore(ytest, ypred)
self.f1 = f1score(ytest, ypred)
self.modelname = "LogisticRegression"
print(f"LR - Accuracy: {self.accuracy:.4f}, F1: {self.f1:.4f}")
self.next(self.compare)
@step
def trainrf(self):
from sklearn.ensemble import RandomForestClassifier
from sklearn.metrics import accuracyscore, f1score
import numpy as np
Xtrain = np.array(self.Xtrain)
ytrain = np.array(self.ytrain)
Xtest = np.array(self.Xtest)
ytest = np.array(self.ytest)
model = RandomForestClassifier(
nestimators=100, randomstate=self.randomstate
)
model.fit(Xtrain, ytrain)
ypred = model.predict(Xtest)
self.accuracy = accuracyscore(ytest, ypred)
self.f1 = f1score(ytest, ypred)
self.modelname = "RandomForest"
print(f"RF - Accuracy: {self.accuracy:.4f}, F1: {self.f1:.4f}")
self.next(self.compare)
@step
def trainsvm(self):
from sklearn.svm import SVC
from sklearn.metrics import accuracyscore, f1score
import numpy as np
Xtrain = np.array(self.Xtrain)
ytrain = np.array(self.ytrain)
Xtest = np.array(self.Xtest)
ytest = np.array(self.ytest)
model = SVC(kernel='rbf', randomstate=self.randomstate)
model.fit(Xtrain, ytrain)
ypred = model.predict(Xtest)
self.accuracy = accuracyscore(ytest, ypred)
self.f1 = f1score(ytest, ypred)
self.modelname = "SVM"
print(f"SVM - Accuracy: {self.accuracy:.4f}, F1: {self.f1:.4f}")
self.next(self.compare)
@step
def compare(self, inputs):
results = []
for inp in inputs:
results.append({
'model': inp.modelname,
'accuracy': inp.accuracy,
'f1': inp.f1
})
results.sort(key=lambda x: x['f1'], reverse=True)
print("\n=== Model Comparison ===")
for r in results:
print(f" {r['model']}: Accuracy={r['accuracy']:.4f}, F1={r['f1']:.4f}")
best = results[0]
self.bestmodelname = best['model']
self.bestaccuracy = best['accuracy']
self.bestf1 = best['f1']
self.allresults = results
print(f"\nBest Model: {self.bestmodelname}")
self.next(self.end)
@step
def end(self):
print("\n=== ML Pipeline Complete ===")
print(f"Selected model: {self.bestmodelname}")
print(f"Accuracy: {self.bestaccuracy:.4f}")
print(f"F1 Score: {self.bestf1:.4f}")
if name == 'main':
MLPipelineFlow()
Run the pipeline:
python mlpipeline.py run
With custom parameters:
python mlpipeline.py run --testsize 0.3 --randomstate 123
Deployment with Argo Workflows
Metaflow supports deployment to Argo Workflows for production scheduling:
python mlpipeline.py argo-workflows create
python mlpipeline.py argo-workflows trigger
Visualization with Metaflow Cards
Metaflow Cards let you create automatic visualizations:
from metaflow import FlowSpec, step, card, current
from metaflow.cards import Markdown, Table, Image
class CardFlow(FlowSpec):
@card
@step
def start(self):
self.metrics = {
'accuracy': 0.95,
'precision': 0.93,
'recall': 0.97
}
current.card.append(Markdown("# Training Results"))
current.card.append(
Table([
['Metric', 'Value'],
['Accuracy', '0.95'],
['Precision', '0.93'],
['Recall', '0.97']
])
)
self.next(self.end)
@step
def end(self):
print("Done!")
if name == 'main':
CardFlow()
View the card after running:
python cardflow.py card view start
Best Practices
1. Design Modular Steps
Each step should perform one specific task. This makes debugging easier and enables per-step retries.
@step
def loaddata(self):
...
self.next(self.validatedata)
@step
def validatedata(self):
...
self.next(self.preprocess)
2. Use Artifacts Wisely
Save only the data that's needed. Avoid storing large objects that won't be accessed in subsequent steps.
@step
def train(self):
self.metrics = {'accuracy': 0.95}
self.next(self.end)
3. Leverage Parameters
Make hyperparameters and configurations into Parameters so experiments are easily reproducible.
4. Use Branching for Experiments
Compare multiple approaches in parallel using branching, then select the best one in the join step.
5. Add Retry for Failure-Prone Steps
Steps that access external APIs or process large data should have retry decorators.
@retry(times=3)
@step
def fetchdata(self):
...
6. Use Namespaces for Organization
METAFLOWUSER=production python flow.py run
METAFLOW
USER=experiment python flow.py run
7. Monitor with Client API
Create monitoring scripts to track pipeline status:
from metaflow import Flow
flow = Flow('MLPipelineFlow')
for run in list(flow.runs())[:5]:
status = "Success" if run.successful else "Failed"
print(f"Run {run.id}: {status} ({run.created_at})")
Conclusion
Metaflow is an MLOps framework designed to simplify data science workflows from experimentation to production. With features like automatic versioning, parallel branching, dependency management, and the Client API, Metaflow enables data scientists to focus on modeling without being burdened by infrastructure complexity.
Key advantages of Metaflow:
- Ease of use: Intuitive and Pythonic API
- Reproducibility: Every experiment is automatically tracked
- Scalability: Seamless transition from laptop to cloud
- Flexibility: Supports multiple cloud providers and orchestrators
- Monitoring: Client API and Cards for observability
For next steps, you can explore Metaflow's integration with AWS Step Functions, Kubernetes, or Apache Airflow for more complex production deployments. The official documentation at docs.metaflow.org provides comprehensive guides for every feature.