Apache Airflow
Apache Airflow adalah platform orkestrasi workflow paling populer di data engineering. Airflow mendefinisikan pipeline sebagai DAG (Directed Acyclic Graph) — kumpulan task dengan dependencies yang jelas.
Konsep Utama
- DAG — Workflow keseluruhan. Mendefinisikan task dan urutannya.
- Task — Unit kerja individual (extract, transform, load).
- Operator — Template untuk task (PythonOperator, BashOperator, dll).
- Schedule — Kapan DAG berjalan (cron expression).
- XCom — Mekanisme passing data antar task.
Menulis DAG Pertama
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
# Default config untuk semua task
default_args = {
"owner": "data-team",
"retries": 2,
"retry_delay": timedelta(minutes=5),
"email_on_failure": True,
"email": ["[email protected]"],
}
# Definisi DAG
with DAG(
dag_id="daily_sales_pipeline",
default_args=default_args,
description="Extract sales data, transform, load to warehouse",
schedule="0 6 * * *", # Setiap hari jam 6 pagi
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["sales", "daily"],
) as dag:
def extract_sales(**context):
execution_date = context["ds"] # YYYY-MM-DD
# ... logic extract
return {"row_count": 1500}
def transform_sales(**context):
ti = context["ti"]
extract_result = ti.xcom_pull(task_ids="extract")
# ... logic transform
def load_to_warehouse(**context):
# ... logic load
pass
extract = PythonOperator(
task_id="extract",
python_callable=extract_sales,
)
transform = PythonOperator(
task_id="transform",
python_callable=transform_sales,
)
load = PythonOperator(
task_id="load",
python_callable=load_to_warehouse,
)
notify = BashOperator(
task_id="notify_slack",
bash_command='curl -X POST "$SLACK_WEBHOOK" -d '{"text":"Pipeline done"}'',
)
# Dependency chain
extract >> transform >> load >> notify
Task Dependencies
# Linear
extract >> transform >> load
# Parallel kemudian converge
[extract_api, extract_db] >> transform >> load
# Branching
extract >> [transform_a, transform_b]
transform_a >> load_a
transform_b >> load_b
Airflow UI
- DAGs view — Daftar semua DAG, status on/off, last run
- Graph view — Visualisasi task dependencies
- Tree view — History run per tanggal
- Log — Output setiap task run untuk debugging