Tutorial Lengkap Prefect: Modern Workflow Orchestration untuk 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

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

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

workpool:

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

workpool:

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:

  • API Pythonic: Native Python decorators
  • Observability: UI dan logging built-in
  • Resilience: Automatic retries dan caching
  • Fleksibilitas: Jalankan dimana saja
  • Integrasi: AWS, GCP, databases, Slack
  • 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

    Artikel Terkait

    Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines

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

    Tutorial Flyte: Workflow Orchestration untuk Machine Learning dan Data Engineering

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

    Tutorial Dagster: Orkestrasi Data dengan Software-Defined Assets

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

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

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