Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines

# Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines Apache Airflow adalah platform open-source untuk membuat, menjadwalkan, dan memonitor workflows secara programatik. Awal...

By Ruby Abdullah · · tutorial
Apache AirflowData PipelineETLMLOpsPythonOrchestration

Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines

Apache Airflow adalah platform open-source untuk membuat, menjadwalkan, dan memonitor workflows secara programatik. Awalnya dikembangkan oleh Airbnb, Airflow telah menjadi standar industri untuk orkestrasi data pipelines dan ML workflows yang kompleks.

Mengapa Airflow?

Keunggulan Airflow:
  • Python-based: Define workflows sebagai code (DAGs)
  • Scalable: Distributed execution dengan Celery/Kubernetes
  • Extensible: Rich ecosystem operators dan hooks
  • Visual UI: Monitor dan manage workflows dengan mudah
  • Active community: Ecosystem besar dan support

Use Cases:
  • ETL/ELT pipelines
  • ML training workflows
  • Data warehouse loading
  • Report generation
  • Infrastructure automation

Instalasi

1. Local Installation

# Buat virtual environment

python -m venv airflowvenv

source airflowvenv/bin/activate

Set Airflow home

export AIRFLOWHOME=~/airflow

Install Airflow

pip install apache-airflow

Initialize database

airflow db init

Buat admin user

airflow users create \

--username admin \

--firstname Admin \

--lastname User \

--role Admin \

--email admin@example.com \

--password admin

Start webserver (di satu terminal)

airflow webserver --port 8080

Start scheduler (di terminal lain)

airflow scheduler

2. Docker Installation

# Download docker-compose file

curl -LfO 'https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml'

Buat directories

mkdir -p ./dags ./logs ./plugins ./config

echo -e "AIRFLOWUID=$(id -u)" > .env

Initialize dan start

docker compose up airflow-init

docker compose up -d

3. Install dengan Extras

# Dengan specific providers

pip install apache-airflow[postgres,google,amazon,slack]

Common extras

pip install apache-airflow[celery] # Celery executor

pip install apache-airflow[kubernetes] # Kubernetes executor

pip install apache-airflow[pandas] # Pandas support

Konsep Dasar

1. DAG (Directed Acyclic Graph)

from datetime import datetime, timedelta

from airflow import DAG

from airflow.operators.python import PythonOperator

from airflow.operators.bash import BashOperator

Default arguments

defaultargs = {

'owner': 'datateam',

'dependsonpast': False,

'email': ['alerts@example.com'],

'emailonfailure': True,

'emailonretry': False,

'retries': 3,

'retrydelay': timedelta(minutes=5),

}

Define DAG

dag = DAG(

'myfirstdag',

defaultargs=defaultargs,

description='DAG tutorial sederhana',

scheduleinterval=timedelta(days=1), # atau '@daily'

startdate=datetime(2024, 1, 1),

catchup=False,

tags=['tutorial'],

)

Define tasks

def helloworld():

print("Hello, World!")

return "success"

task1 = PythonOperator(

taskid='hellotask',

pythoncallable=helloworld,

dag=dag,

)

task2 = BashOperator(

taskid='bashtask',

bashcommand='echo "Hello from Bash"',

dag=dag,

)

Set dependencies

task1 >> task2

from datetime import datetime

from airflow.decorators import dag, task

@dag(

scheduleinterval='@daily',

startdate=datetime(2024, 1, 1),

catchup=False,

tags=['taskflow'],

)

def mytaskflowdag():

@task()

def extract():

"""Extract data dari source"""

data = {"orders": [1, 2, 3, 4, 5]}

return data

@task()

def transform(data: dict):

"""Transform data"""

total = sum(data["orders"])

return {"total": total}

@task()

def load(result: dict):

"""Load data ke destination"""

print(f"Total orders: {result['total']}")

# Define flow

data = extract()

result = transform(data)

load(result)

Instantiate DAG

dag = mytaskflowdag()

3. Schedule Intervals

from airflow import DAG

from datetime import datetime

Preset schedules

dag1 = DAG('dailydag', scheduleinterval='@daily') # Midnight

dag2 = DAG('hourlydag', scheduleinterval='@hourly') # Setiap jam

dag3 = DAG('weeklydag', scheduleinterval='@weekly') # Minggu midnight

dag4 = DAG('monthlydag', scheduleinterval='@monthly') # Hari pertama bulan

Cron expressions

dag5 = DAG('crondag', scheduleinterval='0 6 ') # Jam 6 pagi daily

dag6 = DAG('crondag2', scheduleinterval='0 /2 ') # Setiap 2 jam

dag7 = DAG('crondag3', scheduleinterval='0 9 1-5') # Jam 9 weekdays

timedelta

from datetime import timedelta

dag8 = DAG('timedeltadag', scheduleinterval=timedelta(hours=6))

Operators

1. Python Operator

from airflow.operators.python import PythonOperator, BranchPythonOperator

from airflow.decorators import task

Basic Python Operator

def myfunction(param1, param2):

print(f"Params: {param1}, {param2}")

return param1 + param2

pythontask = PythonOperator(

taskid='pythontask',

pythoncallable=myfunction,

opkwargs={'param1': 10, 'param2': 20},

dag=dag,

)

Dengan XCom

def pushfunction(context):

context['ti'].xcompush(key='mykey', value='myvalue')

def pullfunction(context):

value = context['ti'].xcompull(key='mykey', taskids='pushtask')

print(f"Pulled value: {value}")

Branch Operator

def branchfunction(context):

if context['ds'] == '2024-01-01':

return 'newyeartask'

return 'regulartask'

branchtask = BranchPythonOperator(

taskid='branchtask',

pythoncallable=branchfunction,

dag=dag,

)

2. Bash Operator

from airflow.operators.bash import BashOperator

Simple command

bashtask = BashOperator(

taskid='bashtask',

bashcommand='echo "Hello from Bash"',

dag=dag,

)

Dengan environment variables

bashenvtask = BashOperator(

taskid='bashenvtask',

bashcommand='echo $MYVAR',

env={'MYVAR': 'Hello'},

dag=dag,

)

Run script

bashscripttask = BashOperator(

taskid='bashscripttask',

bashcommand='/path/to/script.sh ',

dag=dag,

)

Dengan templates

bashtemplatetask = BashOperator(

taskid='bashtemplatetask',

bashcommand='echo "Execution date: {{ ds }}"',

dag=dag,

)

3. SQL Operators

from airflow.providers.postgres.operators.postgres import PostgresOperator

from airflow.providers.mysql.operators.mysql import MySqlOperator

PostgreSQL

postgrestask = PostgresOperator(

taskid='postgrestask',

postgresconnid='mypostgres',

sql="""

INSERT INTO mytable (date, value)

VALUES ('{{ ds }}', 100);

""",

dag=dag,

)

MySQL

mysqltask = MySqlOperator(

taskid='mysqltask',

mysqlconnid='mymysql',

sql='SELECT FROM users WHERE createdat >= "{{ ds }}"',

dag=dag,

)

SQL dari file

sqlfiletask = PostgresOperator(

taskid='sqlfiletask',

postgresconnid='mypostgres',

sql='sql/myquery.sql',

dag=dag,

)

4. Sensors

from airflow.sensors.filesystem import FileSensor

from airflow.sensors.externaltask import ExternalTaskSensor

from airflow.providers.http.sensors.http import HttpSensor

File sensor

filesensor = FileSensor(

taskid='waitforfile',

filepath='/data/input/{{ ds }}.csv',

pokeinterval=60, # Check setiap 60 detik

timeout=3600, # Timeout setelah 1 jam

mode='reschedule', # Free up worker slot between checks

dag=dag,

)

External task sensor

externalsensor = ExternalTaskSensor(

taskid='waitforupstream',

externaldagid='upstreamdag',

externaltaskid='finaltask',

timeout=3600,

dag=dag,

)

HTTP sensor

httpsensor = HttpSensor(

taskid='waitforapi',

httpconnid='myapi',

endpoint='/health',

responsecheck=lambda response: response.statuscode == 200,

pokeinterval=30,

timeout=600,

dag=dag,

)

Task Dependencies

1. Basic Dependencies

# Linear dependency

task1 >> task2 >> task3

Sama dengan

task1.setdownstream(task2)

task2.setdownstream(task3)

Multiple downstream

task1 >> [task2, task3]

Multiple upstream

[task1, task2] >> task3

Complex dependencies

task1 >> task2

task1 >> task3

[task2, task3] >> task4

2. Task Groups

from airflow.utils.taskgroup import TaskGroup

from airflow.decorators import dag, task

@dag(scheduleinterval='@daily', startdate=datetime(2024, 1, 1))

def taskgroupdag():

@task

def start():

return "Starting"

with TaskGroup(groupid='processing') as processinggroup:

@task

def processa():

return "Process A"

@task

def processb():

return "Process B"

@task

def combine(a, b):

return f"{a} + {b}"

a = processa()

b = processb()

combine(a, b)

@task

def end():

return "Done"

start() >> processinggroup >> end()

dag = taskgroupdag()

3. Dynamic Task Mapping

from airflow.decorators import dag, task

@dag(scheduleinterval='@daily', startdate=datetime(2024, 1, 1))

def dynamicdag():

@task

def getfiles():

return ['file1.csv', 'file2.csv', 'file3.csv']

@task

def processfile(filename: str):

print(f"Processing {filename}")

return f"Processed {filename}"

@task

def combine(results: list):

print(f"Combined: {results}")

files = getfiles()

processed = processfile.expand(filename=files) # Dynamic mapping

combine(processed)

dag = dynamicdag()

XCom (Cross-Communication)

1. Basic XCom

from airflow.decorators import dag, task

@dag(scheduleinterval='@daily', startdate=datetime(2024, 1, 1))

def xcomdag():

@task

def produce():

# Return value otomatis di-push ke XCom

return {"key": "value", "count": 100}

@task

def consume(data: dict):

# Parameter otomatis di-pull dari XCom

print(f"Received: {data}")

data = produce()

consume(data)

dag = xcomdag()

2. Manual XCom

def pushtask(context):

context['ti'].xcompush(key='customkey', value='customvalue')

def pulltask(context):

# Pull by key

value = context['ti'].xcompull(key='customkey', taskids='pushtask')

# Pull return value

returnvalue = context['ti'].xcompull(taskids='pushtask')

print(f"Key value: {value}, Return value: {returnvalue}")

Connections dan Variables

1. Connections

from airflow.hooks.base import BaseHook

from airflow.providers.postgres.hooks.postgres import PostgresHook

Menggunakan hook

def useconnection(*context):

# Get connection

conn = BaseHook.getconnection('myconnection')

print(f"Host: {conn.host}, Login: {conn.login}")

# Postgres hook

pghook = PostgresHook(postgresconnid='mypostgres')

records = pghook.getrecords("SELECT FROM users")

return records

2. Variables

from airflow.models import Variable

Set variable (via UI atau CLI)

airflow variables set myvar myvalue

Get variable

def usevariable(context):

# Simple get

value = Variable.get('myvar')

# Dengan default

value = Variable.get('myvar', defaultvar='default')

# JSON variable

jsonvar = Variable.get('myjsonvar', deserializejson=True)

return value

Di templates

{{ var.value.myvar }}

{{ var.json.myjsonvar }}

Contoh ETL Pipeline

from datetime import datetime

from airflow.decorators import dag, task

from airflow.providers.postgres.hooks.postgres import PostgresHook

import pandas as pd

@dag(

scheduleinterval='@daily',

startdate=datetime(2024, 1, 1),

catchup=False,

tags=['etl'],

)

def etlpipeline():

@task()

def extract():

"""Extract data dari source database"""

hook = PostgresHook(postgresconnid='sourcedb')

sql = """

SELECT id, name, email, createdat

FROM users

WHERE createdat >= '{{ ds }}'

"""

df = hook.getpandasdf(sql)

return df.todict()

@task()

def transform(data: dict):

"""Transform data"""

df = pd.DataFrame(data)

# Clean data

df['email'] = df['email'].str.lower()

df['name'] = df['name'].str.title()

# Add derived columns

df['domain'] = df['email'].str.split('@').str[1]

df['processedat'] = datetime.now().isoformat()

return df.todict()

@task()

def load(data: dict):

"""Load data ke destination"""

df = pd.DataFrame(data)

hook = PostgresHook(postgresconnid='destdb')

# Insert data

for , row in df.iterrows():

hook.run("""

INSERT INTO processedusers (id, name, email, domain, processedat)

VALUES (%s, %s, %s, %s, %s)

ON CONFLICT (id) DO UPDATE SET

name = EXCLUDED.name,

email = EXCLUDED.email,

domain = EXCLUDED.domain,

processedat = EXCLUDED.processedat

""", parameters=(row['id'], row['name'], row['email'],

row['domain'], row['processedat']))

return len(df)

@task()

def notify(count: int):

"""Kirim notification"""

print(f"ETL complete. Processed {count} records.")

# Define flow

rawdata = extract()

transformeddata = transform(rawdata)

recordcount = load(transformeddata)

notify(recordcount)

dag = etlpipeline()

Contoh ML Pipeline

from datetime import datetime

from airflow.decorators import dag, task

from airflow.providers.amazon.aws.hooks.s3 import S3Hook

@dag(

scheduleinterval='@weekly',

startdate=datetime(2024, 1, 1),

catchup=False,

tags=['ml'],

)

def mltrainingpipeline():

@task()

def fetchdata():

"""Fetch training data dari S3"""

s3hook = S3Hook(awsconnid='awsdefault')

data = s3hook.readkey(

key='data/training/features.csv',

bucketname='my-ml-bucket'

)

return data

@task()

def preprocess(data: str):

"""Preprocess data"""

import pandas as pd

from io import StringIO

df = pd.readcsv(StringIO(data))

# Preprocessing steps

df = df.dropna()

df = df.dropduplicates()

# Feature engineering

df['featureratio'] = df['featurea'] / df['featureb']

return df.todict()

@task()

def trainmodel(data: dict):

"""Train ML model"""

import pandas as pd

from sklearn.ensemble import RandomForestClassifier

from sklearn.modelselection import traintestsplit

import joblib

df = pd.DataFrame(data)

X = df.drop('target', axis=1)

y = df['target']

Xtrain, Xtest, ytrain, ytest = traintestsplit(

X, y, testsize=0.2, randomstate=42

)

model = RandomForestClassifier(nestimators=100)

model.fit(Xtrain, ytrain)

# Simpan model

modelpath = '/tmp/model.joblib'

joblib.dump(model, modelpath)

# Hitung metrics

accuracy = model.score(Xtest, ytest)

return {

'modelpath': modelpath,

'accuracy': accuracy

}

@task()

def evaluatemodel(metrics: dict):

"""Evaluate dan decide deployment"""

if metrics['accuracy'] >= 0.85:

return 'deploy'

return 'skip'

@task()

def deploymodel(metrics: dict, decision: str):

"""Deploy model ke production"""

if decision == 'deploy':

s3hook = S3Hook(awsconnid='awsdefault')

s3hook.loadfile(

filename=metrics['modelpath'],

key=f'models/production/model{{ ds }}.joblib',

bucketname='my-ml-bucket',

replace=True

)

print(f"Model deployed dengan accuracy: {metrics['accuracy']}")

else:

print("Model tidak di-deploy - accuracy di bawah threshold")

# Define flow

rawdata = fetchdata()

processeddata = preprocess(rawdata)

metrics = trainmodel(processeddata)

decision = evaluatemodel(metrics)

deploymodel(metrics, decision)

dag = mltrainingpipeline()

Testing DAGs

# testdag.py

import pytest

from airflow.models import DagBag

from datetime import datetime

@pytest.fixture

def dagbag():

return DagBag()

def testdagloaded(dagbag):

"""Test DAG loaded tanpa error"""

dag = dagbag.getdag('mydag')

assert dag is not None

assert len(dagbag.importerrors) == 0

def testdagtaskcount(dagbag):

"""Test DAG punya jumlah task yang diharapkan"""

dag = dagbag.getdag('mydag')

assert len(dag.tasks) == 4

def testtaskdependencies(dagbag):

"""Test task dependencies"""

dag = dagbag.getdag('mydag')

task1 = dag.gettask('task1')

task2 = dag.gettask('task2')

assert task2.taskid in [t.taskid for t in task1.downstreamlist]

Test individual task

def testpythontask():

"""Test Python callable"""

from dags.mydag import myfunction

result = myfunction(10, 20)

assert result == 30

# Run tests

pytest tests/testdag.py -v

Best Practices

1. DAG Design

# Good practices

from airflow.decorators import dag, task

from datetime import datetime

@dag(

scheduleinterval='@daily',

startdate=datetime(2024, 1, 1),

catchup=False, # Hindari backfill kecuali diperlukan

maxactiveruns=1, # Cegah concurrent runs

tags=['production', 'etl'],

docmd="""

## ETL Pipeline

DAG ini memproses data sales harian.

""",

)

def productiondag():

@task(

retries=3,

retrydelay=timedelta(minutes=5),

executiontimeout=timedelta(hours=1),

)

def extract():

pass

@task(

pool='dbpool', # Gunakan pools untuk limit concurrency

priorityweight=10, # Higher priority

)

def load():

pass

dag = productiondag()

2. Idempotency

@task()

def idempotentload(data: dict, context):

"""Load data secara idempotent"""

executiondate = context['ds']

# Hapus data existing untuk tanggal ini dulu

hook.run(f"DELETE FROM table WHERE date = '{executiondate}'")

# Insert data baru

hook.insertrows('table', data)

3. Resource Management

# airflow.cfg atau environment

Buat pools di UI: Admin -> Pools

@task(pool='databasepool') # Limit concurrent DB connections

def dbtask():

pass

@task(pool='apipool', poolslots=2) # Gunakan 2 slots

def apitask():

pass

Kesimpulan

Apache Airflow adalah standar industri untuk workflow orchestration dengan:

  • Python DAGs: Define workflows sebagai code
  • Rich Operators: SQL, Bash, Python, Cloud providers
  • Scalable: Celery, Kubernetes executors
  • Visual UI: Monitor dan manage workflows
  • Extensible: Custom operators dan hooks
  • Key takeaways:

    • Gunakan TaskFlow API untuk code yang lebih clean
    • Design tasks yang idempotent
    • Gunakan pools untuk manage resources
    • Test DAGs sebelum deployment
    • Monitor dan alert pada failures

    Artikel Terkait

    Tutorial Lengkap Prefect: Modern Workflow Orchestration untuk ML

    Tutorial Lengkap Prefect: Modern Workflow Orchestration untuk ML Prefect adalah platform workflow orchestration modern y...

    Mage: Bikin Data Pipeline Modern yang Rasanya Kayak Ngoding di Notebook

    Mage: Bikin Data Pipeline Modern yang Rasanya Kayak Ngoding di Notebook Halo temen-temen, ketemu lagi sama aku Ruby Abdu...

    Tutorial Dagster: Orkestrasi Data dengan Software-Defined Assets

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

    ZenML: Framework Pipeline MLOps yang Modular dan Cloud-Agnostic

    ZenML: Framework Pipeline MLOps yang Modular dan Cloud-Agnostic Pendahuluan Membangun model machine learning yang akurat...