Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Dynamic task mapping

Airflow & DAGs · Scheduling Deep Dive

Dynamic task mapping

Mediumairflow-47
dynamic-task-mappingexpandtaskflowscalability

Question

What is dynamic task mapping, and when do you use it instead of generating tasks in a loop?

Solution

Dynamic task mapping generates a variable number of task instances at runtime based on the actual output of an upstream task. It replaces the fragile pattern of creating tasks inside static Python loops during DAG parsing.

Parse-time loops versus runtime expansion

In older Airflow DAGs, developers often looped over a list to create tasks:

# Anti-pattern: static parse-time generation
files = ["orders.csv", "users.csv", "events.csv"]
for f in files:
    ProcessOperator(task_id=f"process_{f}", file_name=f)

This required the list of files to be known statically when the DAG file was parsed by the scheduler. Querying an S3 bucket or database at the top level of the DAG file to populate that list hammered the external service every 30 seconds and slowed down the scheduler.

Using expand() and partial()

Dynamic task mapping uses the .expand() method to create tasks dynamically at execution time based on XCom output:

from airflow.decorators import task, dag

@dag(schedule="@daily", start_date=datetime(2025, 1, 1))
def process_batch():
    @task
    def list_files() -> list[str]:
        # Discovers files in S3 at runtime
        return ["s3://b/f1.parquet", "s3://b/f2.parquet", "s3://b/f3.parquet"]

    @task
    def process_file(bucket_name: str, file_path: str):
        # Runs in parallel for each item in list_files
        print(f"Reading {file_path} from {bucket_name}")

    files = list_files()
    # partial sets constant args; expand maps over the runtime list
    process_file.partial(bucket_name="my-bucket").expand(file_path=files)

Operational advantages

Dynamic task mapping provides isolated execution and resilience:

  • Each mapped instance is a discrete task instance with its own state, start time, logs, and retry counters.
  • If file 2 fails due to a corrupt record, only file 2 retries; files 1 and 3 succeed and complete downstream.
  • Airflow provides a safety configuration, max_map_length (defaulting to 1024), which blocks upstream tasks from generating millions of mapped instances that could overwhelm the metadata database.
PreviousNext