Orchestration at Scale — Data Engineering

Orchestration at Scale Setelah Airflow, ada generasi baru orchestrator yang mengatasi limitasi Airflow: Prefect dan Dagster. Keduanya menawarkan developer…

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

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

AspekAirflowPrefectDagster
Mental modelTasks & DAGsTasks & flowsData assets
TestingSulit (butuh env)Mudah (Python biasa)Mudah (Python biasa)
DynamicTerbatasSangat flexibleFlexible
UIFunctionalModernExcellent
MaturityPaling matureGrowingGrowing
CommunityTerbesarMediumMedium

Yang akan kamu pelajari