Dynamic Task Mapping in Apache Airflow: Replace Static Loops with Runtime Parallelism
Airflow 2.3+ dynamic task mapping lets you create task instances at runtime from an upstream list. Example shows S3 file processing with .expand(), plus XCom size and fan‑out guardrails.
02 Nov 2025, 14:09 UTC

The problem with static DAGs
You have a pipeline that processes files dropped into an S3 bucket. The number of files changes every run—sometimes ten, sometimes ten thousand. In classic Airflow you would write a Python loop at DAG parse time to create a task per file, but that only works when the file list is known before the scheduler reads the DAG file. If the list comes from an upstream task (like an S3 listing), the loop cannot see it. The result is either a fragile dynamic DAG generation script or a single monolithic task that processes everything sequentially.
Thesis: dynamic task mapping moves the fan‑out to runtime
Since Airflow 2.3, the TaskFlow API provides .expand() and .partial() so a downstream task can be instantiated once per element of a list produced by an upstream task. The scheduler creates the mapped task instances after the upstream task finishes, so the fan‑out size is decided at runtime. Each instance runs independently, appears as an expandable group in the Grid UI, and can be retried without re‑executing its siblings.
How dynamic mapping works under the hood
- Upstream task returns a Python list (or any iterable) via XCom.
- Downstream task is defined with the
@taskdecorator and called with.expand(arg=upstream_output). - The scheduler reads the XCom value, creates one task instance per list element, and assigns each a
map_index. - Concurrency is controlled by the usual pool limits and the DAG‑level
max_active_tis_per_dagsetting, preventing a massive list from overwhelming workers.
Because the mapped input travels through XCom, the list must stay small enough for the metadata database. A common engineering decision is to pass lightweight references (S3 keys, database IDs) rather than full payloads, and to cap or batch the list size in the upstream task.
Worked example: processing a variable set of S3 files
The following DAG runs on Airflow 2.7+ (syntax is stable across 2.3+ but always verify against your deployed minor version). It lists objects under a prefix, then maps a processing task over each key.
from airflow.decorators import dag, task
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from pendulum import datetime
@dag(
schedule=None,
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["demo", "dynamic-mapping"],
max_active_tis_per_dag=50, # guardrail for fan‑out
)
def process_s3_files():
@task
def list_keys(bucket: str, prefix: str) -> list[str]:
hook = S3Hook(aws_conn_id="aws_default")
keys = hook.list_keys(bucket_name=bucket, prefix=prefix)
# Optional: cap to avoid accidental huge fan‑out
return keys[:500] if keys else []
@task(retries=2)
def process_file(key: str) -> str:
# Replace with actual processing logic
hook = S3Hook(aws_conn_id="aws_default")
content = hook.read_key(key, bucket_name="my-bucket")
# ... transform, load, etc.
return f"processed {key}"
keys = list_keys(bucket="my-bucket", prefix="incoming/")
# Each key becomes a separate mapped task instance
process_file.expand(key=keys)
process_s3_files()
In the Grid view you will see a single process_file task group that expands to show one row per key. If one file fails, you can clear and retry just that instance; the others remain green.
Trade‑offs and guardrails
- XCom size limits: The metadata database (PostgreSQL, MySQL) imposes a practical ceiling on the list size. Returning thousands of keys is usually fine; returning large objects will bloat the DB and slow scheduling.
- Massive fan‑out: Even with
max_active_tis_per_dag, thousands of mapped instances increase scheduler latency and UI rendering time. Use pools on the mapped task (pool="s3-processing") and consider batching the upstream list into chunks. - Downstream aggregation: Collecting results from all mapped instances requires a task that consumes the mapped output (e.g.,
@task def aggregate(results: list[str]): ...called with.expand(results=process_file)). Trigger rules become important—all_successis the default, but you may wantall_doneif partial results are acceptable. - Version drift: The
expand_kwargspattern and mapped operator behavior have evolved across 2.x minor releases. Always test the exact syntax in your target Airflow version before promoting to production.
Actionable next steps
- Identify a DAG that currently uses a parse‑time loop over a runtime‑determined list.
- Prototype the upstream task to return a list of lightweight identifiers.
- Replace the loop with
.expand()on the downstream task, adding a pool andmax_active_tis_per_dagguardrail. - Run the DAG in a staging environment (the official docker‑compose quickstart works well) and verify in the Grid UI that the mapped task expands correctly.
- Inject a failure in one mapped instance (e.g., corrupt a test file) and confirm that retrying it does not re‑run successful siblings.
- Load‑test with a near‑maximum expected list size to observe scheduler latency and metadata DB growth.
Dynamic task mapping eliminates the brittle boundary between DAG parse time and runtime data. By moving the fan‑out into the scheduler, you gain parallelism, observability, and per‑instance retry semantics—without writing custom DAG generation code.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.