Complete Prefect Tutorial: Modern Workflow Orchestration for ML

# Tutorial Lengkap Prefect: Modern Workflow Orchestration untuk ML Prefect adalah platform workflow orchestration modern yang dirancang untuk data dan ML pipelines. Library ini menyediakan cara Pytho...

By Ruby Abdullah · · tutorial
PrefectWorkflow OrchestrationMLOpsData PipelinePythonETL

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

Use Cases:
  • 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

workpool:

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

workpool:

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:

  • Pythonic API: Native Python decorators
  • Observability: Built-in UI and logging
  • Resilience: Automatic retries and caching
  • Flexibility: Run anywhere
  • Integrations: AWS, GCP, databases, Slack
  • 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

    Related Articles

    Complete Apache Airflow Tutorial: Workflow Orchestration for Data Pipelines

    Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines Apache Airflow adalah platform open-source ...

    Flyte Tutorial: Workflow Orchestration for Machine Learning and Data Engineering

    Tutorial Flyte: Workflow Orchestration untuk Machine Learning dan Data Engineering Flyte adalah platform workflow orches...

    Dagster Tutorial: Data Orchestration with Software-Defined Assets

    Dagster: Orkestrasi Data Modern dengan Software-Defined Assets Dagster adalah orkestrator data yang menyusun pipeline be...

    Apache Kafka for Real-Time ML Tutorial: Streaming Data Pipeline

    Tutorial 13: Apache Kafka untuk Pipeline ML Real-Time Daftar Isi Pendahuluan Prasyarat Memahami Apache Kafka [Menyiapkan...