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
ingest transfersCDC lag checkingest KYC vendorAPI pulldbt buildstaging → marts + testsreconcile vs ledgertotals must matchpublish dashboardsand reverse ETLalert owneron failure or lateness
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