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.