Part 5 · 1 chapters · ~8 min
Orchestration: Airflow and Dagster
DAGs of tasks, schedules and data-aware triggers, Airflow operators and sensors, Dagster software-defined assets, retries, backfills and catchup, freshness SLAs, idempotent tasks, and how orchestration differs from durable workflow engines (Workflows course).
7
Pipelines as code
code
# Dagster: assets declare what they produce and depend on
from dagster import asset, AssetCheckResult, asset_check
@asset
def raw_transfers(): ... # loaded by CDC or a sync job
@asset(deps=[raw_transfers])
def fct_daily_volume(): run_dbt(["build", "--select", "fct_daily_volume"])
@asset_check(asset=fct_daily_volume)
def matches_ledger() -> AssetCheckResult:
diff = warehouse_total() - ledger_total() # Ledgers course: the source of truth
return AssetCheckResult(passed=diff == 0, metadata={"diff_kobo": diff})
# Airflow equivalent: a DAG with schedule="@daily", catchup=True for backfills, tasks chained with >>Orchestrators versus workflow engines: Airflow and Dagster schedule batch data jobs measured in minutes to hours; Temporal-style engines run durable business processes with per-entity state (Workflows P6). They solve different problems.
AN ORCHESTRATED DAILY PIPELINE
dependencies, retries and freshness in one graph
swipe the figure sideways, or tap expand for full screen
1/4
a graph of tasks
An orchestrator runs tasks in dependency order on a schedule or when upstream data arrives.
tasks with dependenciesAirflow DAGs, Dagster assets