skip to content

In Airflow, what does calling .expand() on a task produce when the DAG runs?

level: middleimportance: should knowfreq 52%

answer

  1. one declared task, many run-time instances
  2. the width is a property of the run
  3. constants and mapped args go different places
  4. two mapped arguments multiply
  5. zero elements means nothing to run

basics

~20 s

It produces one mapped task instance per element of the expanded input, created at run time once the upstream value is known. The DAG shows a single task node; the run shows N instances, each identified by its map_index.

solid answer

~50 s

Dynamic task mapping lets a task fan out over a collection whose size the DAG file cannot know. `process.expand(key=list_keys())` declares one task; when the run reaches it, the scheduler resolves the upstream XCom, sees N elements, and creates N mapped task instances numbered by `map_index` from 0. Constant arguments go in `.partial(...)`, because `expand` only accepts the ones being mapped: `process.partial(bucket="raw").expand(key=list_keys())`. Expanding over two arguments produces the **cross product**, so build the pairs upstream and use `expand_kwargs` with a list of dicts if you want element-wise pairing. Each instance queues, retries and can be cleared independently, and a downstream non-mapped task receives the list of all results as a natural reduce step. An empty input list skips the mapped task rather than failing it, and a configurable cap limits how many instances one task may create.

code

python · 14 lines
python
@task
def list_keys(bucket: str) -> list[str]:
    return s3_list(bucket)

@task
def process(bucket: str, key: str) -> int:
    return load_one(bucket, key)

@task
def summarise(counts: list[int]) -> None:
    print(sum(counts))

totals = process.partial(bucket="raw-zone").expand(key=list_keys("raw-zone"))
summarise(totals)   # receives the list of every instance's return value

go deeper

for a junior

Recognise the syntax and be able to say that one declared task turns into many instances at run time, one per element of the expanded collection.

for a middle

Explain the split between partial and expand, that multiple expanded arguments produce a cross product, and that instances are indexed by map_index and retry independently.

for a senior

Reason about the operational side: concurrency caps on the fan-out, clearing a single map_index, the skip-on-empty behaviour, and batching an input that could explode.

for a principal

Weigh run-time mapping against parse-time DAG generation for the platform, and set guardrails — maximum widths, pool isolation — so one bad upstream day cannot starve every other DAG on the cluster.

## The problem it solves Before dynamic task mapping, fanning out over a collection meant generating tasks in a Python loop at parse time — which requires knowing the collection while the DAG file is being parsed. That forces the file to query an API or list a bucket on every parse (slow, fragile) and makes the graph shape shift between parses. Dynamic task mapping, added in Airflow 2.3, moves the fan-out to run time. ```python @task def list_keys(bucket: str) -> list[str]: return s3_list(bucket) @task def process(bucket: str, key: str) -> int: return load_one(bucket, key) process.partial(bucket="raw-zone").expand(key=list_keys("raw-zone")) ``` ## What exists at parse time versus run time At parse time there is exactly **one** task in the DAG — a *mapped task*. The UI graph shows one node, annotated with the mapped-instance count once a run exists. When the DAG run reaches that task, the scheduler resolves the `expand` argument (an XComArg pointing at `list_keys`'s output, or a literal list), counts the elements, and expands the mapped task into N **mapped task instances**. Each carries a `map_index` starting at 0; an unmapped instance is denoted by map_index -1. That timing is the whole point and the thing to say out loud: the number of parallel branches is a property of the run, not of the DAG definition, so two runs of the same DAG can have different widths. ## partial versus expand `.expand()` takes only the arguments you want mapped over; everything else must be supplied through `.partial()`. Passing an argument the function does not accept, or omitting a required constant, fails when the DAG is parsed. Classic operators map the same way, with `partial` called on the class: ```python BashOperator.partial(task_id="echo").expand(bash_command=["echo a", "echo b"]) ``` Mapping over multiple keyword arguments yields the **cross product**: `.expand(a=[1, 2], b=["x", "y"])` creates four instances. This is a frequent source of surprise for people expecting pairwise behaviour. For pairing, build the combinations upstream and use `expand_kwargs`, which takes a list of dicts, one dict of arguments per instance: ```python @task def build_args() -> list[dict]: return [{"src": "a", "dst": "A"}, {"src": "b", "dst": "B"}] copy_file.expand_kwargs(build_args()) ``` ## Behaviour of the instances Each mapped instance is a normal task instance in almost every respect: it queues independently, occupies a pool slot, retries on its own schedule, and can be cleared individually from the UI to re-run just the one that failed. How many run at once is bounded by the usual concurrency controls — executor parallelism, the DAG's `max_active_tasks`, pool slots — plus `max_active_tis_per_dag` if you set it on the task. Downstream, a non-mapped task that consumes the mapped task's output receives a **list** of all the instances' return values, which gives you a natural reduce step: ```python totals = process.partial(bucket="raw").expand(key=list_keys("raw")) summarise(totals) # summarise receives [int, int, ...] ``` A downstream task can also itself be mapped over the upstream mapped output, chaining fan-outs. ## Edge cases interviewers probe **Empty list.** If the expanded input resolves to zero elements, there is nothing to run and Airflow marks the mapped task as skipped. Downstream tasks then see a skipped upstream, which under the default trigger rule propagates the skip — a legitimate outcome, but it surprises teams who expected a success with no work. **Unbounded fan-out.** The expanded value comes from data, so a bad upstream day can try to create an enormous number of instances. A configurable maximum map length exists so that exceeding it fails the task rather than melting the scheduler. Treat that as a safety net, not a design: cap or batch the input yourself when the collection could be huge, for example by chunking keys into groups and mapping over the groups. **Not everything is mappable.** The expanded input must be a list-like or dict-like value resolvable at run time, and you still cannot use its length in the DAG file — `len()` of an XComArg is not available at parse time. **Observability.** Logs are per mapped instance, so debugging means picking the failing `map_index`; the grid view groups them under the single task node. Clearing the parent mapped task re-expands from the upstream value, which may produce a different N if the upstream is re-run.

  • What does `.expand(a=[1, 2], b=["x", "y"])` create, and how would you get pairwise behaviour instead?
    Four mapped instances — the cross product of the two lists. For element-wise pairing, build the pairs in an upstream task that returns a list of dicts, then use `expand_kwargs` so each dict supplies one instance's arguments. Doing the pairing upstream also keeps the DAG file free of logic that depends on run-time data.
  • One of fifty mapped instances failed. How do you re-run just that one?
    Clear that single mapped task instance by its map_index from the grid view, or with `airflow tasks clear` scoped to the task and map index; it is scheduled independently and its retries are its own. Clearing the whole mapped task instead re-expands it from the upstream value, which re-runs everything and may even produce a different instance count.
  • How do you consume the results of a mapped task in one downstream step?
    Pass the mapped task's output to an ordinary, unmapped task: it receives a list containing every instance's return value, which is the standard reduce step. Keep those return values small — they are XComs — so return counts, keys or paths rather than the data each instance processed.

saying these in an interview costs you the question

  • Says the fan-out width is fixed when the DAG file is parsed
  • Expects two expanded arguments to be zipped rather than multiplied
  • Puts constant arguments inside expand instead of partial
  • Believes a failed instance forces re-running the whole fan-out
  • Maps over an unbounded list with no cap or batching

context