Complete Prefect Tutorial: Modern Workflow Orchestration for ML
Prefect is a modern workflow orchestration platform designed for data and ML pipelines. It provides a Pythonic way to build, schedule, and monitor complex workflows with automatic retries, caching, and observability.
Why Prefect?
Prefect Advantages:- Pythonic: Native Python decorators and functions
- Observable: Built-in UI and monitoring
- Resilient: Automatic retries and error handling
- Flexible: Run anywhere - local, cloud, Kubernetes
- Modern: Async support and dynamic workflows
- ML pipeline orchestration
- Data engineering workflows
- ETL processes
- Scheduled data processing
- Model training pipelines
Installation
pip install prefect
Verify installation
prefect version
Start Prefect server (local)
prefect server start
Or use Prefect Cloud
prefect cloud login
Quick Start
1. Basic Flow
from prefect import flow, task
@task
def extractdata():
return [1, 2, 3, 4, 5]
@task
def transformdata(data):
return [x 2 for x in data]
@task
def loaddata(data):
print(f"Loading data: {data}")
return len(data)
@flow(name="ETL Pipeline")
def etlpipeline():
rawdata = extractdata()
transformed = transformdata(rawdata)
result = loaddata(transformed)
return result
Run the flow
if name == "main":
etlpipeline()
2. With Parameters
from prefect import flow, task
@task
def fetchdata(url: str):
import requests
response = requests.get(url)
return response.json()
@task
def processdata(data: dict, threshold: float):
return {k: v for k, v in data.items() if v > threshold}
@flow(name="Parameterized Flow")
def datapipeline(url: str, threshold: float = 0.5):
data = fetchdata(url)
processed = processdata(data, threshold)
return processed
Run with parameters
if name == "main":
result = datapipeline(
url="https://api.example.com/data",
threshold=0.7
)
3. ML Training Flow
from prefect import flow, task
import pandas as pd
from sklearn.modelselection import traintestsplit
from sklearn.ensemble import RandomForestClassifier
from sklearn.metrics import accuracyscore
@task
def loaddataset(path: str):
return pd.readcsv(path)
@task
def preprocess(df: pd.DataFrame):
df = df.dropna()
X = df.drop("target", axis=1)
y = df["target"]
return traintestsplit(X, y, testsize=0.2)
@task
def trainmodel(Xtrain, ytrain, nestimators: int):
model = RandomForestClassifier(nestimators=nestimators)
model.fit(Xtrain, ytrain)
return model
@task
def evaluatemodel(model, Xtest, ytest):
predictions = model.predict(Xtest)
accuracy = accuracyscore(ytest, predictions)
return accuracy
@flow(name="ML Training Pipeline")
def trainingpipeline(datapath: str, nestimators: int = 100):
df = loaddataset(datapath)
Xtrain, Xtest, ytrain, ytest = preprocess(df)
model = trainmodel(Xtrain, ytrain, nestimators)
accuracy = evaluatemodel(model, Xtest, ytest)
print(f"Model accuracy: {accuracy:.4f}")
return accuracy
if name == "main":
trainingpipeline("data.csv", nestimators=200)
Task Configuration
1. Retries and Timeouts
from prefect import task, flow
from datetime import timedelta
@task(
retries=3,
retrydelayseconds=10,
timeoutseconds=300
)
def unreliabletask():
import random
if random.random() < 0.5:
raise Exception("Random failure")
return "Success"
@task(
retries=5,
retrydelayseconds=[1, 10, 60, 300, 600] # Exponential backoff
)
def apicall():
import requests
response = requests.get("https://api.example.com")
response.raiseforstatus()
return response.json()
@flow
def resilientflow():
result = unreliabletask()
data = apicall()
return result, data
2. Caching
from prefect import task, flow
from prefect.tasks import taskinputhash
from datetime import timedelta
@task(
cachekeyfn=taskinputhash,
cacheexpiration=timedelta(hours=1)
)
def expensivecomputation(data: list):
import time
time.sleep(10) # Simulate expensive operation
return sum(data)
@task(
cachekeyfn=taskinputhash,
cacheexpiration=timedelta(days=1)
)
def fetchexternaldata(url: str):
import requests
return requests.get(url).json()
@flow
def cachedflow(data: list):
# Second call with same input uses cache
result1 = expensivecomputation(data)
result2 = expensivecomputation(data) # Cached!
return result1, result2
3. Concurrency
from prefect import task, flow
from prefect.tasks import taskinputhash
@task
def processitem(item: int):
import time
time.sleep(1)
return item 2
@flow
def parallelflow():
items = list(range(10))
# Tasks run in parallel by default
results = []
for item in items:
result = processitem.submit(item) # Submit for parallel execution
results.append(result)
# Wait for all results
return [r.result() for r in results]
Using map for parallel processing
@flow
def mapflow():
items = list(range(10))
results = processitem.map(items) # Parallel map
return results
Flow Configuration
1. Flow Settings
from prefect import flow, task
from prefect.logging import getrunlogger
@flow(
name="Configured Flow",
description="A flow with custom configuration",
version="1.0.0",
retries=2,
retrydelayseconds=60,
timeoutseconds=3600,
logprints=True
)
def configuredflow(param: str):
logger = getrunlogger()
logger.info(f"Running with param: {param}")
print("This will also be logged")
return f"Result: {param}"
2. Subflows
from prefect import flow, task
@task
def extract():
return [1, 2, 3]
@task
def transform(data):
return [x 2 for x in data]
@flow
def extractflow():
return extract()
@flow
def transformflow(data):
return transform(data)
@flow(name="Parent Flow")
def parentflow():
# Subflows run as nested flows
rawdata = extractflow()
transformed = transformflow(rawdata)
return transformed
3. Dynamic Flows
from prefect import flow, task
@task
def getdatasources():
return ["sourcea", "sourceb", "sourcec"]
@task
def processsource(source: str):
return f"Processed {source}"
@flow
def dynamicflow():
# Dynamically determine tasks at runtime
sources = getdatasources()
results = []
for source in sources:
result = processsource.submit(source)
results.append(result)
return [r.result() for r in results]
Scheduling and Deployments
1. Create Deployment
from prefect import flow
from prefect.deployments import Deployment
from prefect.server.schemas.schedules import CronSchedule
@flow
def scheduledflow():
print("Running scheduled task")
return "Complete"
Create deployment
deployment = Deployment.buildfromflow(
flow=scheduledflow,
name="daily-run",
schedule=CronSchedule(cron="0 9 "), # 9 AM daily
workqueuename="default"
)
deployment.apply()
2. Using prefect.yaml
# prefect.yaml
name: my-project
prefect-version: 2.0.0
deployments:
- name: ml-training
entrypoint: flows/training.py:trainingpipeline
work
pool:
name: default-agent-pool
schedule:
cron: "0 0 " # Daily at midnight
parameters:
datapath: "s3://bucket/data.csv"
nestimators: 100
- name: data-processing
entrypoint: flows/etl.py:etlpipeline
work
pool:
name: default-agent-pool
schedule:
interval: 3600 # Every hour
# Deploy
prefect deploy --all
3. Work Pools
# Create work pool
prefect work-pool create my-pool --type process
Start worker
prefect worker start --pool my-pool
from prefect import flow
from prefect.deployments import Deployment
@flow
def myflow():
return "Done"
deployment = Deployment.buildfromflow(
flow=myflow,
name="my-deployment",
workpoolname="my-pool"
)
State Management
1. Task States
from prefect import task, flow
from prefect.states import Completed, Failed
@task
def conditionaltask(value: int):
if value < 0:
return Failed(message="Negative value not allowed")
return Completed(data=value 2)
@flow
def stateflow():
result = conditionaltask(5)
print(f"State: {result.type}")
return result
2. Flow States
from prefect import flow, task
from prefect.states import Completed, Failed, Cancelled
@task
def checkcondition():
return True
@task
def maintask():
return "Result"
@flow
def controlledflow():
if not checkcondition():
return Cancelled(message="Condition not met")
result = maintask()
return Completed(data=result)
3. State Handlers
from prefect import flow, task
from prefect.logging import getrunlogger
def onfailure(task, taskrun, state):
logger = getrunlogger()
logger.error(f"Task {task.name} failed: {state.message}")
# Send notification, etc.
@task(onfailure=[onfailure])
def riskytask():
raise Exception("Something went wrong")
@flow
def monitoredflow():
try:
riskytask()
except Exception:
pass
Artifacts and Results
1. Create Artifacts
from prefect import flow, task
from prefect.artifacts import createmarkdownartifact, createtableartifact
@task
def generatereport(data: dict):
markdown = f"""
# Analysis Report
## Summary
- Total records: {data['total']}
- Processed: {data['processed']}
- Failed: {data['failed']}
"""
createmarkdownartifact(
key="analysis-report",
markdown=markdown,
description="Daily analysis report"
)
@task
def createresultstable(results: list):
createtableartifact(
key="results-table",
table=results,
description="Processing results"
)
@flow
def reportingflow():
data = {"total": 100, "processed": 95, "failed": 5}
results = [
{"id": 1, "status": "success"},
{"id": 2, "status": "success"},
{"id": 3, "status": "failed"}
]
generatereport(data)
createresultstable(results)
2. Result Persistence
from prefect import flow, task
from prefect.filesystems import LocalFileSystem
from prefect.results import ResultRecordMetadata
Configure result storage
localstorage = LocalFileSystem(basepath="./results")
@task(resultstorage=localstorage, persistresult=True)
def saveresult(data):
return data
@flow(resultstorage=localstorage, persistresult=True)
def persistentflow():
result = saveresult({"key": "value"})
return result
Integrations
1. With AWS
from prefect import flow, task
from prefectaws import S3Bucket
from prefectaws.credentials import AwsCredentials
@task
def readfroms3(bucketname: str, key: str):
s3bucket = S3Bucket.load("my-s3-bucket")
content = s3bucket.readpath(key)
return content
@task
def writetos3(bucketname: str, key: str, data: bytes):
s3bucket = S3Bucket.load("my-s3-bucket")
s3bucket.writepath(key, data)
@flow
def s3flow():
data = readfroms3("my-bucket", "input/data.csv")
# Process data
writetos3("my-bucket", "output/result.csv", data)
2. With Databases
from prefect import flow, task
from prefectsqlalchemy import SqlAlchemyConnector
@task
def querydatabase(query: str):
connector = SqlAlchemyConnector.load("my-database")
with connector.getconnection() as conn:
result = conn.execute(query)
return result.fetchall()
@task
def insertdata(table: str, data: list):
connector = SqlAlchemyConnector.load("my-database")
with connector.getconnection() as conn:
conn.execute(f"INSERT INTO {table} VALUES (?)", data)
@flow
def databaseflow():
records = querydatabase("SELECT * FROM users")
return records
3. With Slack Notifications
from prefect import flow, task
from prefectslack import SlackWebhook
@task
def sendnotification(message: str):
slack = SlackWebhook.load("my-slack-webhook")
slack.notify(message)
@flow
def notifiedflow():
try:
# Do work
result = "Success"
sendnotification(f"Flow completed: {result}")
except Exception as e:
sendnotification(f"Flow failed: {e}")
raise
Monitoring and Observability
1. Logging
from prefect import flow, task
from prefect.logging import getrunlogger
@task
def loggedtask():
logger = getrunlogger()
logger.debug("Debug message")
logger.info("Info message")
logger.warning("Warning message")
logger.error("Error message")
return "Done"
@flow(logprints=True)
def loggedflow():
print("This will be logged") # Captured with logprints=True
loggedtask()
2. Metrics
from prefect import flow, task
import time
@task
def timedtask():
start = time.time()
# Do work
time.sleep(1)
duration = time.time() - start
return {"duration": duration}
@flow
def metricsflow():
results = []
for i in range(5):
result = timedtask()
results.append(result)
totalduration = sum(r["duration"] for r in results)
return {"totalduration": totalduration}
Best Practices
1. Project Structure
myprefectproject/
flows/
init.py
training.py
etl.py
tasks/
init.py
data.py
ml.py
prefect.yaml
requirements.txt
2. Reusable Tasks
# tasks/data.py
from prefect import task
@task(retries=3, retrydelayseconds=10)
def fetchdata(url: str):
import requests
response = requests.get(url)
response.raiseforstatus()
return response.json()
@task
def validatedata(data: dict, schema: dict):
# Validation logic
return True
flows/etl.py
from prefect import flow
from tasks.data import fetchdata, validatedata
@flow
def etlflow(url: str):
data = fetchdata(url)
isvalid = validatedata(data, {"required": ["id", "name"]})
return data if is_valid else None
Conclusion
Prefect is essential for ML workflow orchestration with:
Key takeaways:
- Use tasks for individual units of work
- Configure retries for unreliable operations
- Enable caching for expensive computations
- Use deployments for scheduled runs
- Leverage artifacts for reporting