Every data platform needs something that runs jobs in the right order, at the right time, and does something sensible when they fail. Apache Airflow is the most widely used open-source tool for that job. It started at Airbnb in 2014, became a top-level Apache project, and reached version 3 in 2025. This article explains the concepts you need, shows a complete daily pipeline, and lists the practices that keep Airflow pipelines reliable.
What Airflow is (and is not)
Airflow is an orchestrator. You describe a workflow in Python as a DAG (directed acyclic graph): a set of tasks and the dependencies between them. Airflow then schedules runs, starts each task when its upstream tasks have succeeded, retries failures, records every run, and shows it all in a web UI.
Airflow is not a data processing engine. Tasks should tell other systems to do heavy work (run SQL in the warehouse, start a Spark job, call dbt) rather than pull millions of rows into Airflow workers. Keeping Airflow as the conductor rather than the orchestra is the single most useful design rule.
Key concepts
| Concept | Meaning |
|---|---|
| DAG | A workflow: tasks plus dependencies, defined in a Python file |
| Task | One unit of work, such as running a query or calling an API |
| Operator | A reusable task template (SQL, Bash, Kubernetes, provider-specific operators for Snowflake, BigQuery, Databricks, dbt, and more) |
| Schedule | When the DAG runs: a cron expression, a preset such as @daily, or data-driven triggers |
| DAG run | One execution of the DAG for a specific data interval (for example, yesterday) |
| Scheduler | The component that decides what should run and queues tasks |
| Executor and workers | Where tasks actually run: locally, on Celery workers, or as Kubernetes pods |
| XCom | Small values passed between tasks, such as a file path or row count (not datasets) |
| Backfill | Running a DAG for past intervals, for example after fixing a bug |
A complete daily pipeline
The DAG below implements the pipeline in the diagram using the TaskFlow API, where each Python function decorated with @task becomes a task and passing one function's result to another creates the dependency. The system-specific calls are left as comments so the structure is easy to see.
from datetime import datetime
from airflow.sdk import dag, task # Airflow 3.x; on 2.x use: from airflow.decorators import dag, task
@dag(
schedule="0 2 * * *", # every day at 02:00
start_date=datetime(2026, 1, 1),
catchup=False,
default_args={"retries": 2},
tags=["warehouse", "daily"],
)
def daily_sales_load():
@task
def extract(ds=None) -> str:
"""Pull yesterday's orders from the source system into object storage."""
path = f"s3://raw-zone/orders/{ds}.parquet"
# ... call the source API or database here ...
return path
@task
def load_raw(path: str) -> int:
"""COPY the file into the warehouse's raw schema; return the row count."""
# ... e.g. warehouse.execute(f"COPY INTO raw.orders FROM '{path}'") ...
return 0
@task
def transform(rows_loaded: int) -> None:
"""Rebuild the modelled tables (for example by running dbt)."""
if rows_loaded == 0:
print("no new rows; skipping transform")
return
# ... subprocess.run(["dbt", "build", "--select", "sales"], check=True) ...
@task
def check_quality() -> None:
"""Fail the run if basic expectations are violated."""
# ... assert no NULL order_ids, revenue within expected range, etc. ...
transform(load_raw(extract())) >> check_quality()
daily_sales_load()A few details worth noticing:
dsis the run's logical date (for example2026-09-27). Using it in file names and queries means each run processes exactly one day, which is what makes re-runs and backfills safe.retries=2handles temporary failures such as a network timeout without waking anyone up.catchup=Falsestops Airflow from automatically running every missed day since the start date when the DAG is first deployed.- The quality check is its own task. If it fails, the run is marked failed and alerts fire, instead of dashboards silently showing bad numbers.
Best practices
- Make every task idempotent. Running a task twice for the same date must give the same result: overwrite the day's partition or use MERGE, never blind appends.
- Keep DAG files light. The scheduler parses DAG files often. Do not query databases or call APIs at the top level of the file; do it inside tasks.
- Push work down. Let the warehouse, Spark or dbt do the processing; Airflow only triggers and monitors.
- Pass references, not data. Use XCom for paths, IDs and counts. Large data belongs in storage.
- Use connections and secrets. Never hard-code credentials; configure them as Airflow connections backed by a secrets manager.
- Alert on failure and on lateness. A pipeline that silently stops is worse than one that fails loudly. Monitor data freshness, not only task status.
- Test DAGs in CI. At minimum, check that every DAG file imports without errors and has no cycles.
Airflow vs the alternatives
| Tool | Best for |
|---|---|
| Apache Airflow | General-purpose orchestration with the largest ecosystem of integrations; managed versions available from major clouds and vendors |
| Dagster | Asset-centric pipelines where you think in terms of tables and their freshness rather than tasks |
| Prefect | Python-first workflows with minimal boilerplate and dynamic flows |
| Built-in schedulers (dbt Cloud, Databricks Workflows, Snowflake tasks) | Simple pipelines that live entirely inside one platform |
If your whole pipeline is "load with a managed connector, then run dbt", a built-in scheduler may be enough. Airflow earns its place once you coordinate several systems, need backfills, or want one place to see every pipeline.
Summary
Airflow turns a set of scripts into a monitored, retryable, backfillable pipeline. Define DAGs in Python, keep tasks idempotent and light, let the warehouse do the heavy work, and treat data-quality checks as first-class tasks. For where orchestration fits in the overall stack, see the modern data engineering architecture guide; for the load pattern this DAG uses, see ETL vs ELT.
Comments
Post a Comment