Orchestration at Scale
Setelah Airflow, ada generasi baru orchestrator yang mengatasi limitasi Airflow: Prefect dan Dagster. Keduanya menawarkan developer experience yang lebih baik dan arsitektur modern.
Limitasi Airflow
- DAG definition yang rigid — perubahan butuh restart scheduler
- Testing lokal sulit — butuh Airflow environment lengkap
- Dynamic pipeline terbatas — DAG harus pre-defined
- UI yang aging — kurang intuitif untuk debugging
- Deployment complexity — banyak komponen (scheduler, webserver, worker, database)
Prefect
Orchestrator modern yang mengutamakan simplicity. Pipeline adalah Python code biasa dengan decorator.
from prefect import flow, task
from prefect.tasks import task_input_hash
from datetime import timedelta
@task(retries=3, cache_key_fn=task_input_hash, cache_expiration=timedelta(hours=1))
def extract(date: str) -> pd.DataFrame:
"""Task biasa — Python function + decorator"""
response = requests.get(f"https://api.example.com/sales?date={date}")
return pd.DataFrame(response.json())
@task(log_prints=True)
def transform(df: pd.DataFrame) -> pd.DataFrame:
cleaned = df.dropna(subset=["amount"])
cleaned["amount"] = cleaned["amount"].astype(float)
print(f"Transformed {len(cleaned)} rows") # Auto-logged
return cleaned
@task
def load(df: pd.DataFrame, table: str):
df.to_sql(table, engine, if_exists="append", index=False)
@flow(name="daily-sales-pipeline")
def sales_pipeline(date: str):
"""Flow = kumpulan tasks"""
raw = extract(date)
clean = transform(raw)
load(clean, "fact_sales")
# Jalankan lokal — tidak perlu server!
sales_pipeline("2024-06-15")
# Atau deploy ke Prefect Cloud
# prefect deployment build sales_pipeline.py:sales_pipeline \
# --name "daily-sales" --cron "0 6 * * *"
Dagster
Orchestrator yang berpusat pada data assets. Bukan "jalankan task A lalu B", tapi "pastikan asset X up-to-date".
from dagster import asset, define_asset_job, Definitions
@asset
def raw_orders() -> pd.DataFrame:
"""Asset: representasi dataset yang dihasilkan"""
return pd.read_csv("s3://data-lake/raw/orders.csv")
@asset
def cleaned_orders(raw_orders: pd.DataFrame) -> pd.DataFrame:
"""Dependency otomatis dari parameter name"""
return raw_orders.dropna().query("amount > 0")
@asset
def revenue_daily(cleaned_orders: pd.DataFrame) -> pd.DataFrame:
return cleaned_orders.groupby("date")["amount"].sum().reset_index()
# Dagster otomatis tahu dependency graph:
# raw_orders → cleaned_orders → revenue_daily
defs = Definitions(
assets=[raw_orders, cleaned_orders, revenue_daily],
jobs=[define_asset_job("daily_refresh", selection="*")]
)
Perbandingan
| Aspek | Airflow | Prefect | Dagster |
|---|---|---|---|
| Mental model | Tasks & DAGs | Tasks & flows | Data assets |
| Testing | Sulit (butuh env) | Mudah (Python biasa) | Mudah (Python biasa) |
| Dynamic | Terbatas | Sangat flexible | Flexible |
| UI | Functional | Modern | Excellent |
| Maturity | Paling mature | Growing | Growing |
| Community | Terbesar | Medium | Medium |