Lesson 15 / 25

Dynamic Task Mapping

Create a variable number of parallel tasks at run time with expand().

One definition, many tasks

Sometimes you only know how many parallel pieces of work you need once the DAG runs, for example "process every file that arrived today". Dynamic task mapping lets a task be expanded over a list so Airflow creates one mapped task instance per item at run time. Use .expand(arg=list_or_xcom) for the varying argument and .partial(...) for fixed arguments. A downstream task can then receive all results as a list. Keep the number of mapped instances reasonable (hundreds, not millions).

Expand over a list

The list_files task returns names; process.expand creates one task per name; summarise receives all return values. I ran it on Airflow 3.0.6 with airflow dags test process_files 2026-10-01; three mapped process instances ran and summarise printed the line below.

from datetime import datetime

from airflow.sdk import dag, task

@dag(schedule="@daily", start_date=datetime(2026, 9, 1), catchup=False)
def process_files():
    @task
    def list_files() -> list[str]:
        return ["a.csv", "b.csv", "c.csv"]

    @task
    def process(name: str) -> int:
        return len(name)            # stand-in for real work

    @task
    def summarise(counts: list[int]) -> None:
        print(f"processed {len(counts)} files")

    summarise(process.expand(name=list_files()))

process_files()

Output:

processed 3 files

Quick check: When is dynamic task mapping useful?

  • When the number of parallel tasks is only known at run time
  • When you want to edit the DAG file while it runs
  • When there are no tasks
  • Never
Answer

When the number of parallel tasks is only known at run time — Mapping generates task instances from data produced during the run.