A task group (@task_group) draws related tasks as one collapsible box and prefixes their IDs (ingest.extract). Unlike Airflow 2 129 's SubDAGs it is purely organizational. The rest of BookNest's DAG:
@dag(
schedule=CronDataIntervalTimetable("0 2 * * *", timezone="UTC"), # ds = day that closed
start_date=datetime(2026, 6, 28, 2, tzinfo=timezone.utc),
end_date=datetime(2026, 6, 30, 2, tzinfo=timezone.utc), # sample orders end on 30 June
catchup=False,
max_active_runs=1, # days load in order, one dbt run at a time
default_args={"retries": 2, "retry_delay": timedelta(minutes=5)},
tags=["booknest"],
)
def booknest_daily():
@task_group
def ingest():
... # extract() and load(), Section 7.5.1
return load(extract())
@task_group
def model():
dbt_build = BashOperator(
task_id="dbt_build",
bash_command=f"{DBT} build --select fct_sales --profiles-dir .")
@task
def check(ds=None) -> str:
lines, gross = p.check_day(ds, booknest_conn())
return f"{lines} lines, gross {gross}"
dbt_build >> check()
@task(outlets=[DAILY_SALES]) # a successful publish updates the asset
def publish(ds=None) -> int:
return p.publish_day(ds, booknest_conn())
ingest() >> model() >> publish()
Business logic lives in an importable module that is unit-testable without Airflow (Unit Testing DAGs with pytest), so the DAG file only wires steps together and parses in a fraction of a second. Retries come from default_args. And the schedule is explicit: CronDataIntervalTimetable makes the 02:00 run cover the previous 24 hours, so ds is the day that just closed, Airflow 2's cron semantics, which Airflow 3 made opt-in.