Apache Airflow Basics — Data Engineering

Apache Airflow Apache Airflow adalah platform orkestrasi workflow paling populer di data engineering. Airflow mendefinisikan pipeline sebagai DAG (Directed…

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

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

Yang akan kamu pelajari