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