kaxil commented on code in PR #74274: URL: https://github.com/apache/airflow/pull/74274#discussion_r4222194711
########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,324 @@ + .. 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. + +.. _loops: + +===== +Loops +===== + +A loop repeats a group of tasks a number of times decided at runtime, +using the results of each iteration. Each iteration can use the previous +iteration's result. Use a loop for work such as refining an answer, following +a hierarchy, or processing successive batches. + +A conditional loop works like a ``do … while`` loop: run the body, then test +whether to run it again. Airflow's ``until`` condition expresses when to stop, +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. + +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 +iterations. An optional ``until`` callable decides when to stop early. + +Airflow evaluates ``until`` in a gate task downstream of the body's terminal task. +It must return a ``bool``: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +Any other return value, including the ``None`` that a forgotten ``return`` produces, will result in +the gate task failing, and the loop not continuing. + +Airflow creates the next iteration when the gate says to continue. +It does not create all ``max_iterations`` up front: +if the loop stops after four iterations, you see those four, without a tail +of skipped iterations. A fixed-count loop runs exactly ``max_iterations`` +iterations. + +This example improves an estimate of the square root of two until the error +is small enough: + +.. 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 +through ``loop.previous``. The gate appears as a task in the loop, named after +the condition function: ``refine.accurate_enough`` in this example. + +Tasks inside a loop body receive ``loop`` as a context parameter, so a parameter with that name in a +loop body must be keyword-only and cannot have a default other than ``None``. + +In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` that hides this one for +the length of the block. Read the Airflow value before the block and use that inside it, for example +``{% set iteration = loop.index %}``. + +When ``until`` has no usable name, the gate is named ``__loop_gate`` instead. +That covers a lambda, a ``functools.partial`` and a callable object. A gate name +that matches a task in the loop body gets a ``__1`` suffix, as with any other +task ID. + +Stopping because ``until`` returned ``True`` means the loop converged. +Reaching ``max_iterations`` while the condition remains ``False`` means the +loop did not converge: the gate task in the final iteration is marked as +failed. Meeting the condition on the last allowed iteration marks the gate task as success. + +For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs Review Comment: The prose says `refine.loop(max_iterations=3)`, but the example below calls `accumulate.loop(max_iterations=3)` (`example_task_loops.py:66`). ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,324 @@ + .. 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. + +.. _loops: + +===== +Loops +===== + +A loop repeats a group of tasks a number of times decided at runtime, +using the results of each iteration. Each iteration can use the previous +iteration's result. Use a loop for work such as refining an answer, following +a hierarchy, or processing successive batches. + +A conditional loop works like a ``do … while`` loop: run the body, then test +whether to run it again. Airflow's ``until`` condition expresses when to stop, +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. + +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 +iterations. An optional ``until`` callable decides when to stop early. + +Airflow evaluates ``until`` in a gate task downstream of the body's terminal task. +It must return a ``bool``: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +Any other return value, including the ``None`` that a forgotten ``return`` produces, will result in +the gate task failing, and the loop not continuing. + +Airflow creates the next iteration when the gate says to continue. +It does not create all ``max_iterations`` up front: +if the loop stops after four iterations, you see those four, without a tail +of skipped iterations. A fixed-count loop runs exactly ``max_iterations`` +iterations. + +This example improves an estimate of the square root of two until the error +is small enough: + +.. 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 +through ``loop.previous``. The gate appears as a task in the loop, named after +the condition function: ``refine.accurate_enough`` in this example. + +Tasks inside a loop body receive ``loop`` as a context parameter, so a parameter with that name in a +loop body must be keyword-only and cannot have a default other than ``None``. + +In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` that hides this one for +the length of the block. Read the Airflow value before the block and use that inside it, for example +``{% set iteration = loop.index %}``. + +When ``until`` has no usable name, the gate is named ``__loop_gate`` instead. +That covers a lambda, a ``functools.partial`` and a callable object. A gate name +that matches a task in the loop body gets a ``__1`` suffix, as with any other +task ID. + +Stopping because ``until`` returned ``True`` means the loop converged. +Reaching ``max_iterations`` while the condition remains ``False`` means the +loop did not converge: the gate task in the final iteration is marked as +failed. Meeting the condition on the last allowed iteration marks the gate task as success. + +For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs +three iterations, carrying results between them. Reaching the cap completes +a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the group. + +.. 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``: + +.. 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 +definition. + +.. list-table:: Loop context + :header-rows: 1 + :widths: 25 75 + + * - Attribute + - Meaning + * - ``loop.index`` + - The current iteration number, starting at zero. + * - ``loop.max_iterations`` + - The configured limit on the number of iterations. + * - ``loop.previous`` + - The previous iteration's result, or ``None`` in iteration 0. + * - ``loop.result`` + - The terminal body task's result (the XCom pushed under the key ``return_value``) for the current iteration, used by the gate. + +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 ``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. Review Comment: Teardowns don't count toward the single terminal task. `get_leaves()` skips them, so for `setup >> work >> [td1, td2]` the gate hangs off `work` next to the teardowns, which `test_shared_terminal_with_two_teardowns_is_one_terminal_definition` asserts. Nothing makes the gate wait for the teardowns, so a teardown failure doesn't stop the loop, and I'd expect the next iteration's setup can start while this iteration's teardowns are still running. Worth stating here, or listing setup/teardown in the body as unsupported. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,324 @@ + .. 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. + +.. _loops: + +===== +Loops +===== + +A loop repeats a group of tasks a number of times decided at runtime, +using the results of each iteration. Each iteration can use the previous +iteration's result. Use a loop for work such as refining an answer, following +a hierarchy, or processing successive batches. + +A conditional loop works like a ``do … while`` loop: run the body, then test +whether to run it again. Airflow's ``until`` condition expresses when to stop, +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. + +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 +iterations. An optional ``until`` callable decides when to stop early. + +Airflow evaluates ``until`` in a gate task downstream of the body's terminal task. +It must return a ``bool``: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +Any other return value, including the ``None`` that a forgotten ``return`` produces, will result in +the gate task failing, and the loop not continuing. + +Airflow creates the next iteration when the gate says to continue. +It does not create all ``max_iterations`` up front: +if the loop stops after four iterations, you see those four, without a tail +of skipped iterations. A fixed-count loop runs exactly ``max_iterations`` +iterations. + +This example improves an estimate of the square root of two until the error +is small enough: + +.. 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 +through ``loop.previous``. The gate appears as a task in the loop, named after +the condition function: ``refine.accurate_enough`` in this example. + +Tasks inside a loop body receive ``loop`` as a context parameter, so a parameter with that name in a +loop body must be keyword-only and cannot have a default other than ``None``. + +In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` that hides this one for +the length of the block. Read the Airflow value before the block and use that inside it, for example +``{% set iteration = loop.index %}``. + +When ``until`` has no usable name, the gate is named ``__loop_gate`` instead. +That covers a lambda, a ``functools.partial`` and a callable object. A gate name +that matches a task in the loop body gets a ``__1`` suffix, as with any other +task ID. + +Stopping because ``until`` returned ``True`` means the loop converged. +Reaching ``max_iterations`` while the condition remains ``False`` means the +loop did not converge: the gate task in the final iteration is marked as +failed. Meeting the condition on the last allowed iteration marks the gate task as success. + +For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs +three iterations, carrying results between them. Reaching the cap completes +a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the group. + +.. 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``: + +.. 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 +definition. + +.. list-table:: Loop context + :header-rows: 1 + :widths: 25 75 + + * - Attribute + - Meaning + * - ``loop.index`` + - The current iteration number, starting at zero. + * - ``loop.max_iterations`` + - The configured limit on the number of iterations. + * - ``loop.previous`` + - The previous iteration's result, or ``None`` in iteration 0. + * - ``loop.result`` + - The terminal body task's result (the XCom pushed under the key ``return_value``) for the current iteration, used by the gate. + +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 ``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`` +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. + +If the body has several branches, finish with a task that combines their +outputs. For example, if two branches end in ``refine_left`` and +``refine_right``, add this inside the task group: + +.. code-block:: python + + @task + def combine(left, right): + return {"left": left, "right": right} + + + combine(refine_left(), refine_right()) + +``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. See :ref:`dag-result`. + +Connect a loop to other tasks +============================= + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + refinement >> finished() + +With the default ``all_success`` trigger rule, ``finished`` waits for successful +loop completion. Other trigger rules behave normally; ``always`` does not wait +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 +collection of items. + +A mapped task can be the body's single terminal task: + +.. code-block:: text + + select items → process each item → gate + +The gate reads the collection of mapped results through ``loop.result``. Add +a combining task before the gate if you want to reduce those results to a +single value or structure. A zero-length expansion skips the mapped task and, +under the gate's ``all_success`` rule, skips the gate; no next iteration is +created. + +In this example, ``choose_items`` returns the values to map over: ``[1, 2]`` in the +first iteration, then each previous result plus one. + +.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py + :start-after: [START mapped_loop] + :end-before: [END mapped_loop] + +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. + +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 +task with ``all_done`` can handle a branch failure and return a result. If that +task succeeds, the gate evaluates its result normally. If the terminal task +fails, the gate's ``all_success`` rule prevents it from running, and +no next iteration is created. You cannot change the gate task's trigger rule. +The gate has exactly one upstream, the body's terminal task, so set +``trigger_rule`` on that task to change when the gate runs. + +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, as described under non-convergence above. + +Skipping the body's terminal task also skips the gate under its +``all_success`` trigger rule, so no next iteration is created. Downstream +tasks with the default ``all_success`` rule are skipped too; give a downstream +task ``trigger_rule="none_failed"`` if it should still run when the loop ends +without running: + +.. code-block:: python + + @task(trigger_rule="none_failed") + def finished(): ... + + + refinement >> finished() + +If some branches may be skipped but the loop should continue, give the final +combining task a trigger rule that permits those skips and have it return the +iteration's result. + +Manually marking a gate successful completes it without evaluating ``until`` +or creating another iteration. This is an explicit override of normal gate +execution. Downstream tasks then run as they would after any successful gate: they cannot tell a +gate that was marked successful by hand from one that stopped because ``until`` returned ``True``. +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. + +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: + +* **Downstream** includes downstream tasks. It starts selected unless you have Review Comment: I pointed you at the wrong label here, sorry. For a task inside a loop the Clear action opens `ClearExecutionDialog`, whose checkbox reads "Clear downstream tasks" (`execution.clearDownstream` in `dag.json`). **Downstream** is the label in the dialog used for tasks outside a loop. Lines 281 and 298 need the same change. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,324 @@ + .. 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. + +.. _loops: + +===== +Loops +===== + +A loop repeats a group of tasks a number of times decided at runtime, +using the results of each iteration. Each iteration can use the previous +iteration's result. Use a loop for work such as refining an answer, following +a hierarchy, or processing successive batches. + +A conditional loop works like a ``do … while`` loop: run the body, then test +whether to run it again. Airflow's ``until`` condition expresses when to stop, +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. + +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 +iterations. An optional ``until`` callable decides when to stop early. + +Airflow evaluates ``until`` in a gate task downstream of the body's terminal task. +It must return a ``bool``: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +Any other return value, including the ``None`` that a forgotten ``return`` produces, will result in +the gate task failing, and the loop not continuing. + +Airflow creates the next iteration when the gate says to continue. +It does not create all ``max_iterations`` up front: +if the loop stops after four iterations, you see those four, without a tail +of skipped iterations. A fixed-count loop runs exactly ``max_iterations`` +iterations. + +This example improves an estimate of the square root of two until the error +is small enough: + +.. 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 +through ``loop.previous``. The gate appears as a task in the loop, named after +the condition function: ``refine.accurate_enough`` in this example. + +Tasks inside a loop body receive ``loop`` as a context parameter, so a parameter with that name in a +loop body must be keyword-only and cannot have a default other than ``None``. + +In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` that hides this one for +the length of the block. Read the Airflow value before the block and use that inside it, for example +``{% set iteration = loop.index %}``. + +When ``until`` has no usable name, the gate is named ``__loop_gate`` instead. +That covers a lambda, a ``functools.partial`` and a callable object. A gate name +that matches a task in the loop body gets a ``__1`` suffix, as with any other +task ID. + +Stopping because ``until`` returned ``True`` means the loop converged. +Reaching ``max_iterations`` while the condition remains ``False`` means the +loop did not converge: the gate task in the final iteration is marked as +failed. Meeting the condition on the last allowed iteration marks the gate task as success. + +For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs +three iterations, carrying results between them. Reaching the cap completes +a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the group. + +.. 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``: + +.. 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 +definition. + +.. list-table:: Loop context + :header-rows: 1 + :widths: 25 75 + + * - Attribute + - Meaning + * - ``loop.index`` + - The current iteration number, starting at zero. + * - ``loop.max_iterations`` + - The configured limit on the number of iterations. + * - ``loop.previous`` + - The previous iteration's result, or ``None`` in iteration 0. + * - ``loop.result`` + - The terminal body task's result (the XCom pushed under the key ``return_value``) for the current iteration, used by the gate. + +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 ``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`` +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. + +If the body has several branches, finish with a task that combines their +outputs. For example, if two branches end in ``refine_left`` and +``refine_right``, add this inside the task group: + +.. code-block:: python + + @task + def combine(left, right): + return {"left": left, "right": right} + + + combine(refine_left(), refine_right()) + +``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. Review Comment: I don't think this recipe can be followed. A task outside the loop can't pull a loop task's XCom: at the top of the stack `producer_contexts` raises "A loop producer requires a consumer inside the loop" (`models/task_coordinates.py:386`), and an `xcom_pull` with no scope comes back as a 400 from `execution_api/routes/xcoms.py`. The next bullet and the 409 message in `core_api/services/public/dag_run.py:308` give the same advice. Could the page show a way that works, for example writing the converged value somewhere durable from inside the loop, or is there a supported way for an outside task to read the final iteration's result? ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,324 @@ + .. 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. + +.. _loops: + +===== +Loops +===== + +A loop repeats a group of tasks a number of times decided at runtime, +using the results of each iteration. Each iteration can use the previous +iteration's result. Use a loop for work such as refining an answer, following +a hierarchy, or processing successive batches. + +A conditional loop works like a ``do … while`` loop: run the body, then test +whether to run it again. Airflow's ``until`` condition expresses when to stop, +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. + +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 +iterations. An optional ``until`` callable decides when to stop early. + +Airflow evaluates ``until`` in a gate task downstream of the body's terminal task. +It must return a ``bool``: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +Any other return value, including the ``None`` that a forgotten ``return`` produces, will result in +the gate task failing, and the loop not continuing. + +Airflow creates the next iteration when the gate says to continue. +It does not create all ``max_iterations`` up front: +if the loop stops after four iterations, you see those four, without a tail +of skipped iterations. A fixed-count loop runs exactly ``max_iterations`` +iterations. + +This example improves an estimate of the square root of two until the error +is small enough: + +.. 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 +through ``loop.previous``. The gate appears as a task in the loop, named after +the condition function: ``refine.accurate_enough`` in this example. + +Tasks inside a loop body receive ``loop`` as a context parameter, so a parameter with that name in a +loop body must be keyword-only and cannot have a default other than ``None``. + +In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` that hides this one for +the length of the block. Read the Airflow value before the block and use that inside it, for example +``{% set iteration = loop.index %}``. + +When ``until`` has no usable name, the gate is named ``__loop_gate`` instead. +That covers a lambda, a ``functools.partial`` and a callable object. A gate name +that matches a task in the loop body gets a ``__1`` suffix, as with any other +task ID. + +Stopping because ``until`` returned ``True`` means the loop converged. +Reaching ``max_iterations`` while the condition remains ``False`` means the +loop did not converge: the gate task in the final iteration is marked as +failed. Meeting the condition on the last allowed iteration marks the gate task as success. + +For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs +three iterations, carrying results between them. Reaching the cap completes +a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the group. + +.. 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``: + +.. 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 +definition. + +.. list-table:: Loop context + :header-rows: 1 + :widths: 25 75 + + * - Attribute + - Meaning + * - ``loop.index`` + - The current iteration number, starting at zero. + * - ``loop.max_iterations`` + - The configured limit on the number of iterations. + * - ``loop.previous`` + - The previous iteration's result, or ``None`` in iteration 0. + * - ``loop.result`` + - The terminal body task's result (the XCom pushed under the key ``return_value``) for the current iteration, used by the gate. + +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 ``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`` +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. + +If the body has several branches, finish with a task that combines their +outputs. For example, if two branches end in ``refine_left`` and +``refine_right``, add this inside the task group: + +.. code-block:: python + + @task + def combine(left, right): + return {"left": left, "right": right} + + + combine(refine_left(), refine_right()) + +``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. See :ref:`dag-result`. + +Connect a loop to other tasks +============================= + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + refinement >> finished() + +With the default ``all_success`` trigger rule, ``finished`` waits for successful +loop completion. Other trigger rules behave normally; ``always`` does not wait +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 +collection of items. + +A mapped task can be the body's single terminal task: + +.. code-block:: text + + select items → process each item → gate + +The gate reads the collection of mapped results through ``loop.result``. Add +a combining task before the gate if you want to reduce those results to a +single value or structure. A zero-length expansion skips the mapped task and, +under the gate's ``all_success`` rule, skips the gate; no next iteration is +created. + +In this example, ``choose_items`` returns the values to map over: ``[1, 2]`` in the +first iteration, then each previous result plus one. + +.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py + :start-after: [START mapped_loop] + :end-before: [END mapped_loop] + +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. + +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 +task with ``all_done`` can handle a branch failure and return a result. If that +task succeeds, the gate evaluates its result normally. If the terminal task +fails, the gate's ``all_success`` rule prevents it from running, and +no next iteration is created. You cannot change the gate task's trigger rule. Review Comment: This isn't quite true on the current code. The gate is created without a `trigger_rule` (`definitions/_internal/loop.py:113`), so it picks one up from `default_args` like any other task in the group, and `test_gate_inherits_body_trigger_rule_default` checks that for `all_done` and `one_success`. With `default_args={"trigger_rule": "all_done"}`, a fixed-count loop whose terminal task fails still runs the gate, and the gate continues because it only checks `loop.index + 1 >= loop.max_iterations`. Either pin the gate's trigger rule or say here that `default_args` reaches it. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,324 @@ + .. 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. + +.. _loops: + +===== +Loops +===== + +A loop repeats a group of tasks a number of times decided at runtime, +using the results of each iteration. Each iteration can use the previous +iteration's result. Use a loop for work such as refining an answer, following +a hierarchy, or processing successive batches. + +A conditional loop works like a ``do … while`` loop: run the body, then test +whether to run it again. Airflow's ``until`` condition expresses when to stop, +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. + +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 +iterations. An optional ``until`` callable decides when to stop early. + +Airflow evaluates ``until`` in a gate task downstream of the body's terminal task. +It must return a ``bool``: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +Any other return value, including the ``None`` that a forgotten ``return`` produces, will result in +the gate task failing, and the loop not continuing. + +Airflow creates the next iteration when the gate says to continue. +It does not create all ``max_iterations`` up front: +if the loop stops after four iterations, you see those four, without a tail +of skipped iterations. A fixed-count loop runs exactly ``max_iterations`` +iterations. + +This example improves an estimate of the square root of two until the error +is small enough: + +.. 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 +through ``loop.previous``. The gate appears as a task in the loop, named after +the condition function: ``refine.accurate_enough`` in this example. + +Tasks inside a loop body receive ``loop`` as a context parameter, so a parameter with that name in a +loop body must be keyword-only and cannot have a default other than ``None``. + +In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` that hides this one for +the length of the block. Read the Airflow value before the block and use that inside it, for example +``{% set iteration = loop.index %}``. + +When ``until`` has no usable name, the gate is named ``__loop_gate`` instead. +That covers a lambda, a ``functools.partial`` and a callable object. A gate name +that matches a task in the loop body gets a ``__1`` suffix, as with any other +task ID. + +Stopping because ``until`` returned ``True`` means the loop converged. +Reaching ``max_iterations`` while the condition remains ``False`` means the +loop did not converge: the gate task in the final iteration is marked as +failed. Meeting the condition on the last allowed iteration marks the gate task as success. + +For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs +three iterations, carrying results between them. Reaching the cap completes +a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the group. + +.. 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``: + +.. 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 +definition. + +.. list-table:: Loop context + :header-rows: 1 + :widths: 25 75 + + * - Attribute + - Meaning + * - ``loop.index`` + - The current iteration number, starting at zero. + * - ``loop.max_iterations`` + - The configured limit on the number of iterations. + * - ``loop.previous`` + - The previous iteration's result, or ``None`` in iteration 0. + * - ``loop.result`` + - The terminal body task's result (the XCom pushed under the key ``return_value``) for the current iteration, used by the gate. + +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 ``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`` +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. + +If the body has several branches, finish with a task that combines their +outputs. For example, if two branches end in ``refine_left`` and +``refine_right``, add this inside the task group: + +.. code-block:: python + + @task + def combine(left, right): + return {"left": left, "right": right} + + + combine(refine_left(), refine_right()) + +``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. See :ref:`dag-result`. + +Connect a loop to other tasks +============================= + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + refinement >> finished() + +With the default ``all_success`` trigger rule, ``finished`` waits for successful +loop completion. Other trigger rules behave normally; ``always`` does not wait +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 +collection of items. + +A mapped task can be the body's single terminal task: + +.. code-block:: text + + select items → process each item → gate + +The gate reads the collection of mapped results through ``loop.result``. Add +a combining task before the gate if you want to reduce those results to a +single value or structure. A zero-length expansion skips the mapped task and, +under the gate's ``all_success`` rule, skips the gate; no next iteration is +created. + +In this example, ``choose_items`` returns the values to map over: ``[1, 2]`` in the +first iteration, then each previous result plus one. + +.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py + :start-after: [START mapped_loop] + :end-before: [END mapped_loop] + +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. + +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 +task with ``all_done`` can handle a branch failure and return a result. If that +task succeeds, the gate evaluates its result normally. If the terminal task +fails, the gate's ``all_success`` rule prevents it from running, and +no next iteration is created. You cannot change the gate task's trigger rule. +The gate has exactly one upstream, the body's terminal task, so set +``trigger_rule`` on that task to change when the gate runs. + +An exception in ``until`` fails the gate task and appears in its logs. If the Review Comment: The gate also takes `retries` from `default_args`. Non-convergence raises `LoopMaxIterationsExceeded`, a plain `AirflowException`, so with `retries=3, retry_delay=5min` the final gate retries three times against the same terminal result and fails about 15 minutes late. Only the non-bool return raises `AirflowFailException` and skips retries. Raising `AirflowFailException` for non-convergence as well, or saying here that the gate retries, would settle it. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
