Loops and mapped tasks

A Dag does not have to be static. How much work it does can depend on data that only exists at runtime: how many files arrived, what the previous step found, or whether an answer is good enough yet. Airflow has two features for this, as well as mechanisms like retries and dynamic Dag generation that fit adjacent use cases.

Choose a mechanism

You want to

Use

How it works

Run the same task once for each item in a collection

Mapped tasks

task.expand(...) creates one task instance per item once the collection is known. The instances can run concurrently.

Repeat a group of tasks, each iteration building on the result of the previous one, until a condition is met

Loops

task_group.loop(...) creates the next iteration of task instances only if the last iteration says another is needed.

Give a failed task another try

Retries, set with retries on the task. See Tasks.

A retry reruns the same task instance. It does not create new task instances, and it does not advance a loop.

Create tasks from something known when the Dag file is parsed

A Python for loop in the Dag file. See Dynamic Dag Generation.

The Dag has the same shape in every run.

Use them together

An iteration of a loop can contain mapped tasks, so each iteration can fan out over a different collection of items. See Mapped tasks inside a loop.

Mapping a whole loop, nesting loops, and placing a loop inside a mapped task group are not supported.

Was this entry helpful?