Complete Apache Airflow Tutorial: Workflow Orchestration for 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

Complete Apache Airflow Tutorial: Workflow Orchestration for Data Pipelines

Apache Airflow is an open-source platform to programmatically author, schedule, and monitor workflows. Originally developed by Airbnb, Airflow has become the industry standard for orchestrating complex data pipelines and ML workflows.

Why Airflow?

Airflow Advantages:
  • Python-based: Define workflows as code (DAGs)
  • Scalable: Distributed execution with Celery/Kubernetes
  • Extensible: Rich ecosystem of operators and hooks
  • Visual UI: Monitor and manage workflows easily
  • Active community: Large ecosystem and support

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

Installation

1. Local Installation

# Create 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

Create admin user

airflow users create \

--username admin \

--firstname Admin \

--lastname User \

--role Admin \

--email admin@example.com \

--password admin

Start webserver (in one terminal)

airflow webserver --port 8080

Start scheduler (in another terminal)

airflow scheduler

2. Docker Installation

# Download docker-compose file

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

Create directories

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

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

Initialize and start

docker compose up airflow-init

docker compose up -d

3. Install with Extras

# With 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

Core Concepts

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='A simple tutorial DAG',

scheduleinterval=timedelta(days=1), # or '@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 from source"""

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

return data

@task()

def transform(data: dict):

"""Transform the data"""

total = sum(data["orders"])

return {"total": total}

@task()

def load(result: dict):

"""Load data to 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') # Every hour

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

dag4 = DAG('monthlydag', scheduleinterval='@monthly') # First day of month

Cron expressions

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

dag6 = DAG('crondag2', scheduleinterval='0 /2 ') # Every 2 hours

dag7 = DAG('crondag3', scheduleinterval='0 9 1-5') # 9 AM 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,

)

With 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,

)

With 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,

)

With 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 from 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 every 60 seconds

timeout=3600, # Timeout after 1 hour

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

Same as

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 is automatically pushed to XCom

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

@task

def consume(data: dict):

# Parameter is automatically pulled from 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 and Variables

1. Connections

from airflow.hooks.base import BaseHook

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

Using 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 or CLI)

airflow variables set myvar myvalue

Get variable

def usevariable(context):

# Simple get

value = Variable.get('myvar')

# With default

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

# JSON variable

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

return value

In templates

{{ var.value.myvar }}

{{ var.json.myjsonvar }}

ETL Pipeline Example

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 from 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 the 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 to 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):

"""Send notification"""

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

# Define flow

rawdata = extract()

transformeddata = transform(rawdata)

recordcount = load(transformeddata)

notify(recordcount)

dag = etlpipeline()

ML Pipeline Example

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 from 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 the 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 the 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)

# Save model

modelpath = '/tmp/model.joblib'

joblib.dump(model, modelpath)

# Calculate metrics

accuracy = model.score(Xtest, ytest)

return {

'modelpath': modelpath,

'accuracy': accuracy

}

@task()

def evaluatemodel(metrics: dict):

"""Evaluate and decide deployment"""

if metrics['accuracy'] >= 0.85:

return 'deploy'

return 'skip'

@task()

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

"""Deploy model to 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 with accuracy: {metrics['accuracy']}")

else:

print("Model not deployed - accuracy below 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 is loaded without errors"""

dag = dagbag.getdag('mydag')

assert dag is not None

assert len(dagbag.importerrors) == 0

def testdagtaskcount(dagbag):

"""Test DAG has expected number of tasks"""

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, # Avoid backfill unless needed

maxactiveruns=1, # Prevent concurrent runs

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

docmd="""

## ETL Pipeline

This DAG processes daily sales data.

""",

)

def productiondag():

@task(

retries=3,

retrydelay=timedelta(minutes=5),

executiontimeout=timedelta(hours=1),

)

def extract():

pass

@task(

pool='dbpool', # Use pools to limit concurrency

priorityweight=10, # Higher priority

)

def load():

pass

dag = productiondag()

2. Idempotency

@task()

def idempotentload(data: dict, context):

"""Load data idempotently"""

executiondate = context['ds']

# Delete existing data for this date first

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

# Insert new data

hook.insertrows('table', data)

3. Resource Management

# airflow.cfg or environment

Create pools in UI: Admin -> Pools

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

def dbtask():

pass

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

def apitask():

pass

Conclusion

Apache Airflow is the industry standard for workflow orchestration with:

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

    • Use TaskFlow API for cleaner code
    • Design idempotent tasks
    • Use pools to manage resources
    • Test DAGs before deployment
    • Monitor and alert on failures

    Related Articles

    Complete Prefect Tutorial: Modern Workflow Orchestration for ML

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

    Mage: Building Modern Data Pipelines That Feel Like Coding in a Notebook

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

    Dagster Tutorial: Data Orchestration with Software-Defined Assets

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

    ZenML: Modular and Cloud-Agnostic MLOps Pipeline Framework

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