In Airflow, what does calling .expand() on a task produce when the DAG runs?
answer
- one declared task, many run-time instances
- the width is a property of the run
- constants and mapped args go different places
- two mapped arguments multiply
- zero elements means nothing to run
basics
~20 sIt 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 sDynamic 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@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 valuego deeper
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.
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.
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.
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