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
- 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
2. TaskFlow API (Recommended)
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 test
dagloaded(dagbag):
"""Test DAG loaded tanpa error"""
dag = dagbag.get
dag('mydag')
assert dag is not None
assert len(dagbag.import
errors) == 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:
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