Apache Airflow is the most-used data orchestrator in 2026. Despite newer alternatives, it remains the default in most enterprises. Knowing it is core to data engineering.
The core concepts
DAG (Directed Acyclic Graph)
A workflow definition. Nodes are tasks; edges define dependencies.
from airflow.decorators import dag, task
from datetime import datetime, timedelta
@dag(
dag_id="daily_customer_pipeline",
start_date=datetime(2026, 1, 1),
schedule="0 3 * * *", # 3am daily
catchup=False,
)
def customer_pipeline():
@task
def extract_orders():
# ...
return order_count
@task
def extract_customers():
# ...
return customer_count
@task
def transform_data(order_count, customer_count):
# ...
return "ok"
@task
def load_to_dashboard(status):
# ...
orders = extract_orders()
customers = extract_customers()
transformed = transform_data(orders, customers)
load_to_dashboard(transformed)
dag = customer_pipeline()
The TaskFlow API (@dag/@task decorators) is the modern Airflow 2.x+ style. Use it.
Task
A unit of work. Operators encapsulate work (PythonOperator, BashOperator, SQLExecuteQueryOperator, etc.).
Schedule
When the DAG runs. Cron expression or interval. schedule="@daily", "0 3 * * *", etc.
Sensor
A task that waits for an external condition (file exists, table updated, API ready).
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
wait_for_file = S3KeySensor(
task_id="wait_for_daily_file",
bucket_name="incoming",
bucket_key="customer_data_{{ ds }}.csv",
poke_interval=60,
timeout=3600,
)
XCom
Inter-task data passing (small values; not for large datasets).
@task
def extract():
return {"row_count": 1000}
@task
def report(stats):
print(stats["row_count"]) # consumes the dict from XCom
DAG design principles
1. Idempotency
Same DAG run twice = same result. Critical for retries and backfills.
Bad:
@task
def add_count():
db.execute("INSERT INTO summary (count) VALUES (...)") # accumulates on retry
Good:
@task(execution_date_in_task=True)
def add_count(execution_date):
db.execute("DELETE FROM summary WHERE date = %s", execution_date)
db.execute("INSERT INTO summary (date, count) VALUES (...)", execution_date)
The DELETE-then-INSERT pattern is the most common idempotency idiom in DE.
2. Atomic tasks
One task = one logical operation. Don't bundle 5 operations in one task.
Bad:
@task
def do_everything():
extract()
transform()
load()
Good:
@task
def extract(): ...
@task
def transform(extraction_result): ...
@task
def load(transformation_result): ...
Atomic tasks let you:
- Retry just the failed part.
- See clear timing per step.
- Understand failures in logs.
3. Configuration as code
DAGs are Python. Use Python features for repetition:
SOURCES = ["stripe", "salesforce", "shopify"]
@task
def extract_from_source(source_name):
return extract(source_name)
@task
def load_source_data(source_name, data):
load(source_name, data)
for src in SOURCES:
data = extract_from_source.override(task_id=f"extract_{src}")(src)
load_source_data.override(task_id=f"load_{src}")(src, data)
Don't copy-paste DAG code for each source.
4. Avoid heavy compute in Airflow workers
Airflow workers are for orchestration, not for processing big data.
Bad:
@task
def process_huge_dataset():
df = pd.read_csv("100GB-file.csv") # OOMs the worker
return df.aggregate(...)
Good:
@task
def trigger_spark_job():
from spark_submit_operator import submit_spark_job
return submit_spark_job(jar_path="...")
Push compute to where it belongs (Spark, dbt, Snowflake). Airflow orchestrates; doesn't process.
5. Use connections and variables for config
Don't hardcode credentials/URLs:
from airflow.providers.postgres.hooks.postgres import PostgresHook
hook = PostgresHook(postgres_conn_id="warehouse_prod")
hook.run("SELECT ...")
warehouse_prod is configured in Airflow's UI/secrets manager. Never in code.
Common DAG patterns
Pattern 1: Extract → load → dbt run
@dag(...)
def daily_pipeline():
extract = PythonOperator(task_id="extract", ...)
load = PythonOperator(task_id="load", ...)
transform = BashOperator(
task_id="dbt_run",
bash_command="cd /dbt && dbt run --target prod",
)
extract >> load >> transform
dag = daily_pipeline()
Standard daily pipeline.
Pattern 2: Source-fanout
SOURCES = ["stripe", "shopify", "salesforce"]
@dag(...)
def fanout():
@task
def extract(source): ...
@task
def load(source, data): ...
for src in SOURCES:
data = extract.override(task_id=f"extract_{src}")(src)
load.override(task_id=f"load_{src}")(src, data)
All sources extract in parallel; loads independently.
Pattern 3: Sensor + downstream
@dag(...)
def file_processing():
wait_for_file = S3KeySensor(
task_id="wait_for_input",
bucket_name="incoming",
bucket_key="data_{{ ds }}.csv",
)
process = PythonOperator(task_id="process", ...)
wait_for_file >> process
DAG waits for the daily file to appear, then processes.
What Airflow does well
- Industry standard. Knowledge transfers across jobs.
- Massive operator library (any source/sink has one).
- Mature scheduling.
- Wide ecosystem (Astronomer, MWAA, Cloud Composer for managed).
What Airflow does poorly
- Verbose. Boilerplate-heavy.
- Schedule timezone confusion (
logical_datevs run time). - Backfills are painful.
- Heavy infrastructure for small teams (need a database, scheduler, workers).
- Pythonic but not Pythonic enough (the Airflow API has its own conventions).
Modern alternatives (Dagster, Prefect — next lesson) address these but have less industry penetration.
Common Airflow mistakes
- Heavy compute in tasks. Push to Spark/dbt/SQL.
- Hardcoded credentials. Use connections.
- Non-idempotent tasks. Retries cause duplicates.
- Long-running tasks (>4h). Hard to debug, expensive to restart.
- Catchup enabled by default. New DAG fires historical runs unexpectedly.
- DAGs that import heavy libraries at top level. Slows scheduler.
Production Airflow deployment
Three options:
- Self-host on Kubernetes — most flexibility, most ops burden.
- Managed services — Astronomer, MWAA (AWS), Cloud Composer (GCP). $$$.
- Docker Compose — for small teams; ok for prototyping/dev.
For startups, MWAA or Cloud Composer covers needs at moderate cost. For enterprise, self-hosted or Astronomer.
Takeaway
Airflow = the orchestrator everyone knows. Use TaskFlow API (decorators). Atomic, idempotent tasks. Don't compute heavy data in workers. Use connections for credentials. Push to Spark/dbt for processing. Despite alternatives, Airflow remains the default; learn it well.