Metaflow Tutorial: Netflix's MLOps Framework for Data Science

# Tutorial Metaflow: Framework MLOps dari Netflix untuk Data Science Metaflow adalah framework open-source yang dikembangkan oleh Netflix untuk membangun dan mengelola proyek data science secara efis...

By Ruby Abdullah · · tutorial
MetaflowMLOpsNetflixPipelinePython

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

@condabase(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

METAFLOWUSER=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.

Related Articles

ClearML Tutorial: Open-Source MLOps Platform for Experiment Tracking and Pipeline Automation

Tutorial ClearML: Platform MLOps Open-Source untuk Experiment Tracking dan Pipeline Automation ClearML adalah platform M...

Kedro Tutorial: Reproducible and Maintainable Data Science Pipelines

Kedro: Pipeline Data Science yang Reproducible dan Mudah Dirawat Sebagian besar proyek data science dimulai dari satu no...

ZenML: Modular and Cloud-Agnostic MLOps Pipeline Framework

ZenML: Framework Pipeline MLOps yang Modular dan Cloud-Agnostic Pendahuluan Membangun model machine learning yang akurat...

Complete Comet ML Tutorial: MLOps Platform for Experiment Tracking and Model Management

Tutorial Lengkap Comet ML: Platform MLOps untuk Experiment Tracking dan Model Management Dalam dunia machine learning mo...