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 Pythonic untuk membangun, menjadwalkan, dan memonitor workflow kompleks dengan automatic retries, caching, dan observability.
Mengapa Prefect?
Keunggulan Prefect:- Pythonic: Decorators dan fungsi Python native
- Observable: UI dan monitoring built-in
- Resilient: Automatic retries dan error handling
- Fleksibel: Jalankan dimana saja - local, cloud, Kubernetes
- Modern: Async support dan dynamic workflows
- ML pipeline orchestration
- Data engineering workflows
- Proses ETL
- Scheduled data processing
- Model training pipelines
Instalasi
pip install prefect
Verify instalasi
prefect version
Start Prefect server (local)
prefect server start
Atau gunakan 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
Jalankan flow
if name == "main":
etlpipeline()
2. Dengan 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
Jalankan dengan parameter
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)
Konfigurasi Task
1. Retries dan 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("Kegagalan random")
return "Sukses"
@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) # Simulasi operasi mahal
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):
# Panggilan kedua dengan input sama menggunakan 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 berjalan paralel secara default
results = []
for item in items:
result = processitem.submit(item) # Submit untuk eksekusi paralel
results.append(result)
# Tunggu semua hasil
return [r.result() for r in results]
Menggunakan map untuk processing paralel
@flow
def mapflow():
items = list(range(10))
results = processitem.map(items) # Parallel map
return results
Konfigurasi Flow
1. Flow Settings
from prefect import flow, task
from prefect.logging import getrunlogger
@flow(
name="Configured Flow",
description="Flow dengan konfigurasi custom",
version="1.0.0",
retries=2,
retrydelayseconds=60,
timeoutseconds=3600,
logprints=True
)
def configuredflow(param: str):
logger = getrunlogger()
logger.info(f"Berjalan dengan param: {param}")
print("Ini juga akan di-log")
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 berjalan sebagai 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"Diproses {source}"
@flow
def dynamicflow():
# Tentukan tasks secara dinamis saat runtime
sources = getdatasources()
results = []
for source in sources:
result = processsource.submit(source)
results.append(result)
return [r.result() for r in results]
Scheduling dan Deployments
1. Buat Deployment
from prefect import flow
from prefect.deployments import Deployment
from prefect.server.schemas.schedules import CronSchedule
@flow
def scheduledflow():
print("Menjalankan task terjadwal")
return "Complete"
Buat deployment
deployment = Deployment.buildfromflow(
flow=scheduledflow,
name="daily-run",
schedule=CronSchedule(cron="0 9 "), # 9 AM setiap hari
workqueuename="default"
)
deployment.apply()
2. Menggunakan 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 " # Setiap hari tengah malam
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 # Setiap jam
# Deploy
prefect deploy --all
3. Work Pools
# Buat 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="Nilai negatif tidak diizinkan")
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="Kondisi tidak terpenuhi")
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} gagal: {state.message}")
# Kirim notifikasi, dll.
@task(onfailure=[onfailure])
def riskytask():
raise Exception("Sesuatu salah")
@flow
def monitoredflow():
try:
riskytask()
except Exception:
pass
Artifacts dan Results
1. Buat Artifacts
from prefect import flow, task
from prefect.artifacts import createmarkdownartifact, createtableartifact
@task
def generatereport(data: dict):
markdown = f"""
# Laporan Analisis
## Ringkasan
- Total records: {data['total']}
- Diproses: {data['processed']}
- Gagal: {data['failed']}
"""
createmarkdownartifact(
key="analysis-report",
markdown=markdown,
description="Laporan analisis harian"
)
@task
def createresultstable(results: list):
createtableartifact(
key="results-table",
table=results,
description="Hasil processing"
)
@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
Konfigurasi 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
Integrasi
1. Dengan 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")
# Proses data
writetos3("my-bucket", "output/result.csv", data)
2. Dengan 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. Dengan Notifikasi Slack
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:
# Lakukan pekerjaan
result = "Sukses"
sendnotification(f"Flow selesai: {result}")
except Exception as e:
sendnotification(f"Flow gagal: {e}")
raise
Monitoring dan Observability
1. Logging
from prefect import flow, task
from prefect.logging import getrunlogger
@task
def loggedtask():
logger = getrunlogger()
logger.debug("Pesan debug")
logger.info("Pesan info")
logger.warning("Pesan warning")
logger.error("Pesan error")
return "Done"
@flow(logprints=True)
def loggedflow():
print("Ini akan di-log") # Ditangkap dengan logprints=True
loggedtask()
2. Metrics
from prefect import flow, task
import time
@task
def timedtask():
start = time.time()
# Lakukan pekerjaan
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. Struktur Project
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):
# Logika validasi
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
Kesimpulan
Prefect adalah essential untuk workflow orchestration ML dengan:
Key takeaways:
- Gunakan tasks untuk unit kerja individual
- Konfigurasi retries untuk operasi tidak reliable
- Enable caching untuk komputasi mahal
- Gunakan deployments untuk scheduled runs
- Manfaatkan artifacts untuk reporting