This is an automated email from the ASF dual-hosted git repository.
ashb pushed a commit to branch task-loops-stack-1
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/task-loops-stack-1 by this
push:
new 6ec5dd6c749 fixup! Add user-facing docs for the new Task Loops feature.
6ec5dd6c749 is described below
commit 6ec5dd6c7495ce1290eba9407121980b2e3bad4d
Author: Ash Berlin-Taylor <[email protected]>
AuthorDate: Tue Oct 6 22:49:28 2026 +0100
fixup! Add user-facing docs for the new Task Loops feature.
Review threads T1.13, T1.17, T1.19, T1.20, T1.25, T1.28, T1.32, T1.33,
T1.35, T1.36, T1.39, T1.40 and T1.3: the page told readers to detect the first
iteration with loop.previous is None, which is also true when the previous
terminal task returned nothing, so the examples now use loop.index == 0. It
also lacked a working partial()/override() example, left a malformed RST list
that swallowed the next heading, used a single-backtick literal, used the
undefined start()/finish() names and [...]
The guide's inline snippets were copies that nothing ran or checked.
Keeping the example Dags beside the guide in the docs tree, and including them
from there, gives the snippets one source from the start. They are not
importable Dags yet because the authoring API arrives later in the series,
which moves them into the example Dags package.
Review thread T1.40. The guide named no view for browsing iterations. The
run-level page for a task group has one Task Instances tab with an Iteration
filter (including All iterations) and an Iteration column, so the text now says
to select the group in a Dag run and use that filter, instead of describing a
view that groups instances by iteration.
Co-Authored-By: Claude <[email protected]>
---
.../dynamic-task-mapping.rst | 2 +-
.../examples/example_task_loops.py | 110 ++++++++++++++++
.../loops-and-mapped-tasks.rst | 2 +-
.../docs/authoring-and-scheduling/loops.rst | 146 +++++++--------------
4 files changed, 162 insertions(+), 98 deletions(-)
diff --git
a/airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rst
b/airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rst
index ef549932c89..972f7440011 100644
--- a/airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rst
+++ b/airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rst
@@ -392,7 +392,7 @@ For example, this code will *not* work:
When code in ``my_task_group`` is executed, ``value`` would still only be a
reference, not the real value, so the ``if not value`` branch will not work as
you likely want. However, if you pass that reference into a task, it will
become resolved when the task is executed, and the three ``my_task`` instances
will therefore receive 1, 2, and 3, respectively.
-It is, therefore, important to remember that, if you intend to perform any
logic on a value passed into a task group function, you must always use a task
to run the logic, such as ``@task.branch`` (or ``BranchPythonOperator``) for
conditions, and task mapping methods for loops.
+It is, therefore, important to remember that, if you intend to perform any
logic on a value passed into a task group function, you must always use a task
to run the logic, such as ``@task.branch`` (or ``BranchPythonOperator``) for
conditions, and ``expand()`` to iterate over values. To repeat a group of
tasks, see :ref:`loops-and-mapped-tasks`.
.. note:: Task-mapping in a mapped task group is not permitted
diff --git
a/airflow-core/docs/authoring-and-scheduling/examples/example_task_loops.py
b/airflow-core/docs/authoring-and-scheduling/examples/example_task_loops.py
new file mode 100644
index 00000000000..89ec4827e7a
--- /dev/null
+++ b/airflow-core/docs/authoring-and-scheduling/examples/example_task_loops.py
@@ -0,0 +1,110 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Repeat task groups with fixed counts, runtime conditions, and mapped
tasks."""
+
+from __future__ import annotations
+
+# [START refine_estimate]
+from airflow.sdk import dag, task, task_group
+
+
+@dag(schedule=None, catchup=False, tags=["example"])
+def refine_estimate():
+ @task_group
+ def refine():
+ @task
+ def improve(*, loop):
+ estimate = 1.0 if loop.index == 0 else loop.previous["estimate"]
+ return (estimate + 2.0 / estimate) / 2.0
+
+ @task
+ def evaluate(estimate):
+ return {"estimate": estimate, "error": abs(estimate * estimate -
2.0)}
+
+ evaluate(improve())
+
+ def accurate_enough(*, loop):
+ return loop.result["error"] < 0.000001
+
+ @task
+ def finished():
+ print("Refinement finished.")
+
+ refinement = refine.loop(max_iterations=10, until=accurate_enough)
+ refinement >> finished()
+
+
+refine_estimate()
+# [END refine_estimate]
+
+
+# [START fixed_loop]
+@dag(schedule=None, catchup=False, tags=["example"])
+def fixed_task_loop():
+ @task_group
+ def accumulate():
+ @task
+ def increment(*, loop):
+ return (0 if loop.index == 0 else loop.previous) + 1
+
+ increment()
+
+ accumulate.loop(max_iterations=3)
+
+
+fixed_task_loop()
+# [END fixed_loop]
+
+
+# [START mapped_loop]
+@dag(schedule=None, catchup=False, tags=["example"])
+def mapped_task_loop():
+ @task_group
+ def process_batch():
+ @task
+ def process(value, *, loop, ti):
+ print(f"Iteration {loop.index}, mapped position {ti.map_index}")
+ return value + loop.index
+
+ process.expand(value=[1, 2])
+
+ def batch_ready(*, loop):
+ return min(loop.result) >= 2
+
+ process_batch.loop(max_iterations=3, until=batch_ready)
+
+
+mapped_task_loop()
+# [END mapped_loop]
+
+
+# [START partial_override_loop]
+@dag(schedule=None, catchup=False, tags=["example"])
+def partial_override_loop():
+ @task_group
+ def accumulate(increment_by):
+ @task
+ def increment(increment_by, *, loop):
+ return (0 if loop.index == 0 else loop.previous) + increment_by
+
+ increment(increment_by)
+
+
accumulate.override(group_id="accumulate_more").partial(increment_by=5).loop(max_iterations=3)
+
+
+partial_override_loop()
+# [END partial_override_loop]
diff --git
a/airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst
b/airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst
index ec6bcf52601..7b8621313ca 100644
--- a/airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst
+++ b/airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst
@@ -55,7 +55,7 @@ Choose a mechanism
Use them together
=================
-A pass of a loop can contain mapped tasks, so each pass can fan out over a
different collection of
+An iteration of a loop can contain mapped tasks, so each iteration can fan out
over a different collection of
items. See :ref:`Mapped tasks inside a loop <loops-mapped-tasks>`.
Mapping a whole loop, nesting loops, and placing a loop inside a mapped task
group are not supported.
diff --git a/airflow-core/docs/authoring-and-scheduling/loops.rst
b/airflow-core/docs/authoring-and-scheduling/loops.rst
index 6fc91b6aef4..092c41c1d0c 100644
--- a/airflow-core/docs/authoring-and-scheduling/loops.rst
+++ b/airflow-core/docs/authoring-and-scheduling/loops.rst
@@ -17,6 +17,7 @@
.. _loops:
+=====
Loops
=====
@@ -31,13 +32,13 @@ so it corresponds to ``do { body } while not
until(result)``. The body runs at
least once. ``max_iterations`` sets an upper limit; the runtime condition
determines how many iterations are needed within that limit.
-Task instances within an iteration can run in parallel.
+Iterations run one after another; task instances within an iteration can run
in parallel.
If you are not sure whether you need a loop or mapped tasks, see
:ref:`loops-and-mapped-tasks`.
Create a loop
--------------
+=============
Define the body with ``@task_group`` and call ``.loop()`` on the decorated
function. Supply a positive integer ``max_iterations`` to limit the number of
@@ -57,39 +58,9 @@ count is reached.
This example improves an estimate of the square root of two until the error
is small enough:
-.. code-block:: python
-
- from airflow.sdk import dag, task, task_group
-
-
- @dag
- def refine_estimate():
- @task_group
- def refine():
- @task
- def improve(*, loop):
- previous = loop.previous
- estimate = 1.0 if previous is None else previous["estimate"]
- return (estimate + 2.0 / estimate) / 2.0
-
- @task
- def evaluate(estimate):
- return {"estimate": estimate, "error": abs(estimate * estimate
- 2.0)}
-
- evaluate(improve())
-
- def accurate_enough(*, loop):
- return loop.result["error"] < 0.000001
-
- @task
- def finished():
- print("Refinement finished.")
-
- refinement = refine.loop(max_iterations=10, until=accurate_enough)
- refinement >> finished()
-
-
- refine_estimate()
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+ :start-after: [START refine_estimate]
+ :end-before: [END refine_estimate]
``evaluate`` returns the result for the iteration. The gate reads it through
``loop.result``. If another iteration runs, ``improve`` reads that same result
@@ -110,30 +81,20 @@ For a fixed-count loop, omit ``until``. The definition
``refine.loop(max_iterati
three iterations, carrying results between them. Reaching the cap completes
a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the
group.
-.. code-block:: python
-
- @dag
- def fixed_task_loop():
- @task_group
- def accumulate():
- @task
- def increment(*, loop):
- previous = loop.previous
- return (0 if previous is None else previous) + 1
-
- increment()
-
- accumulate.loop(max_iterations=3)
-
-
- fixed_task_loop()
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+ :start-after: [START fixed_loop]
+ :end-before: [END fixed_loop]
For a task-group function with arguments, supply them with ``.partial()``
before
calling ``.loop()``. Use ``.override()`` to configure the group, for example to
-give another loop a different ``group_id``.
+give another loop a different ``group_id``:
+
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+ :start-after: [START partial_override_loop]
+ :end-before: [END partial_override_loop]
Read the loop context
----------------------
+=====================
Declare ``loop`` as a keyword-only parameter on a task function. Airflow
supplies it at execution time; leave it out when calling the task in the Dag
@@ -158,14 +119,15 @@ Use ``loop.index`` for the loop iteration and
``ti.map_index`` for a mapped
task instance's position. Each task instance can have several tries, each with
its own ``ti.try_number``.
Pass data between iterations
-----------------------------
+============================
-Between iterations, use ``loop.previous``. Check explicitly for ``None`` to
-handle the first iteration; a previous result of ``0``, ``False``, or an empty
-collection can be valid data.
+Between iterations, use ``loop.previous``. Check ``loop.index == 0`` to handle
+the first iteration. Do not test ``loop.previous is None``: it is also ``None``
+when the previous terminal task returned nothing, and a previous result of
+``0``, ``False``, or an empty collection can be valid data.
The loop body must have exactly one terminal task definition before the gate.
-Its return value (the XCom pushed under the `return_value` key) becomes
``loop.result`` for the gate and ``loop.previous``
+Its return value (the XCom pushed under the ``return_value`` key) becomes
``loop.result`` for the gate and ``loop.previous``
for the next iteration. A mapped terminal task supplies its collection of
results as a sequence (``LazyXComSequence``), the same as other downstream
consumers of mapped tasks.
@@ -185,27 +147,29 @@ outputs. For example, if two branches end in
``refine_left`` and
``combine`` is now the single terminal task. The gate can read each result
through ``loop.result["left"]`` and ``loop.result["right"]``; tasks in the
next iteration use the corresponding keys in ``loop.previous``.
+
Limitations:
* ``include_prior_dates=True`` cannot select a loop iteration from another Dag
run. Push the result to XCom through a task outside the loop if later Dag runs
need to retrieve it with an ordinary XCom pull.
-* The experimental DagRun wait API also requires an outside-loop result task.
It rejects results selected directly from a loop member.
+* The experimental DagRun wait API also requires an outside-loop result task.
It rejects results selected directly from a loop member. See :ref:`dag-result`.
+
Connect a loop to other tasks
-------------------------------
+=============================
Use the object returned by ``.loop()`` in dependencies:
.. code-block:: python
- start() >> refinement >> finish()
+ start() >> refinement >> finished()
With the default ``all_success`` trigger rule, ``finish`` waits for successful
loop completion. Other trigger rules behave normally; ``always`` does not wait
-for the loop.
+for the loop.
.. _loops-mapped-tasks:
Mapped tasks inside a loop
---------------------------
+==========================
A loop iteration can contain mapped tasks. Use mapping to process several
items within the iteration. The next iteration can operate on a different
@@ -223,29 +187,13 @@ single value or structure. A zero-length expansion skips
the mapped task and,
under the gate's default trigger rule, skips the gate; no next iteration is
created.
-.. code-block:: python
-
- @dag(schedule=None, catchup=False, tags=["example"])
- def mapped_task_loop():
- @task_group
- def process_batch():
- @task
- def process(value, *, loop, ti):
- print(f"Iteration {loop.index}, mapped position {ti.map_index}")
- return value + loop.index
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+ :start-after: [START mapped_loop]
+ :end-before: [END mapped_loop]
- process.expand(value=[1, 2])
-
- def batch_ready(*, loop):
- return min(loop.result) >= 2
-
- process_batch.loop(max_iterations=3, until=batch_ready)
-
-
- mapped_task_loop()
-
-The loop iteration and mapped position are separate "coordinates". Mapped
-instance 2 in iteration 0 is distinct from mapped instance 2 in iteration 1.
+The loop iteration and mapped position are separate "coordinates". In this
+example, each iteration has mapped instances 0 and 1, so mapped instance 1 in
+iteration 0 is distinct from mapped instance 1 in iteration 1.
Mapped instances appear within their loop iteration, so you can inspect their
states, tries, and logs separately.
@@ -253,7 +201,7 @@ Mapping a whole loop, nesting loops, and placing a loop
inside a mapped task
group are not supported.
Failures, skips, and retries
-----------------------------
+============================
A retry stays in the same iteration and does not consume another iteration.
Tasks in the body follow normal trigger rules. For example, a final combining
@@ -264,8 +212,8 @@ no next iteration is created.
An exception in ``until`` fails the gate task and appears in its logs. If the
loop reaches its iteration limit without meeting ``until``, the final gate
-fails. Downstream tasks with
-the default ``all_success`` trigger rule will not run.
+fails; downstream tasks with the default ``all_success`` trigger rule will not
+run, as described under non-convergence above.
Skipping the body's terminal task also skips the gate under its default
``all_success`` trigger rule, so no next iteration is created. Downstream
@@ -284,14 +232,20 @@ or creating another iteration. This is an explicit
override of normal gate
execution. Any later iterations retained after a selective clear remain
unchanged.
Iterations and execution history
---------------------------------
+================================
+
+In the Grid, select the loop's task group in a Dag run to open its Task
Instances
+tab. The **Iteration** filter narrows the table to one iteration, or shows
+**All iterations**, and the **Iteration** column shows which iteration each
task
+instance belongs to. The gate's state and logs explain why the loop continued,
+stopped, or failed. Each task's tries and logs remain accessible within its
+iteration.
-The loop view groups tasks and mapped instances by iteration. The gate's state
-and logs explain why the loop continued, stopped, or failed. Each task's tries
-and logs remain accessible within its iteration.
+Iterations cleared by a rerun are kept as history but are not shown in the UI
+in 3.4.0.
Clear tasks inside a loop
---------------------------
+==========================
Use these controls to select how far a clear extends through the loop:
@@ -308,7 +262,7 @@ iterations again. The gate can stop earlier or continue
further, within the
configured limit.
Rerun part of an iteration
-~~~~~~~~~~~~~~~~~~~~~~~~~~~
+---------------------------
Suppose each iteration contains this sequence:
@@ -329,7 +283,7 @@ If the gate now stops at iteration 2, replacement
iterations 3 and 4 are not
created.
Keep later iterations
-~~~~~~~~~~~~~~~~~~~~~
+---------------------
Deselect **Clear later loop iterations** to keep the later iterations.
Rerunning an earlier gate then does not change the loop progression that