Task Groups

Task Groups and DAG Organization

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:

dags/booknest_daily.py (excerpt): schedule, task groups and the asset outletPython
@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()
The graph view of a booknest_daily run: the ingest and model task groups expanded, then publish
The graph view of a booknest_daily run: the ingest and model task groups expanded, then 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.