TJaniF commented on code in PR #74274: URL: https://github.com/apache/airflow/pull/74274#discussion_r4194046852
########## airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst: ########## @@ -0,0 +1,61 @@ + .. 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-and-mapped-tasks: + +====================== +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, and sometimes neither is what you need. + +Choose a mechanism +================== + +.. list-table:: + :header-rows: 1 + :widths: 30 30 40 + + * - You want to + - Use + - How it works + * - Run the same task once for each item in a collection + - :ref:`Mapped tasks <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 pass building on the result of the previous one, until a + condition is met + - :ref:`Loops <loops>` + - ``task_group.loop(...)`` creates the next pass of task instances only if the last pass says Review Comment: ```suggestion * - Repeat a group of tasks, each iteration building on the result of the previous one, until a condition is met - :ref:`Loops <loops>` - ``task_group.loop(...)`` creates the next iteration of task instances only if the last iteration says ``` The loops doc uses "iteration" which I think is easier to understand. ########## airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst: ########## @@ -0,0 +1,61 @@ + .. 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-and-mapped-tasks: + +====================== +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, and sometimes neither is what you need. + +Choose a mechanism +================== + +.. list-table:: + :header-rows: 1 + :widths: 30 30 40 + + * - You want to + - Use + - How it works + * - Run the same task once for each item in a collection + - :ref:`Mapped tasks <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 pass building on the result of the previous one, until a + condition is met + - :ref:`Loops <loops>` + - ``task_group.loop(...)`` creates the next pass of task instances only if the last pass says + another is needed. + * - Give a failed task another attempt + - Retries, set with ``retries`` on the task. See :doc:`/core-concepts/tasks`. + - A retry reruns the same task instance. It does not create new tasks, and it does not advance a Review Comment: ```suggestion - A retry reruns the same task instance. It does not create new task instances, and it does not advance a ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. Review Comment: ```suggestion least once. ``max_iterations`` sets an upper limit; the runtime condition determines how many iterations are needed within that limit. ``` Just a preference, I think "limit" is easier to understand for ESL people. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: Review Comment: ```suggestion Airflow evaluates ``until`` in an automatically generated gate task downstream of the body's terminal task: ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +count is reached. Review Comment: ```suggestion of skipped iterations. For a fixed-count loop, when no ``until`` callable is provided, the gate always returns ``False`` until the count is reached. ``` If I understood it correctly and this is to say there is still a gate task even if I don't provide an "until" callable but it will always return False until max_iterations is reached. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound the number of Review Comment: ```suggestion function. Supply a positive integer ``max_iterations`` to limit the number of ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. Review Comment: ```suggestion - The terminal body task's result (the XCom pushed under the key ``return_value``) for the current iteration, used by the gate. ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + Review Comment: ```suggestion 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``. .. code-block:: python @task_group def accumulate(): @task def increment(increment_by, loop): previous = loop.previous return (0 if previous is None else previous) + increment_by increment() accumulate.override( group_id="accumulate_more" ).partial( increment_by=5 ).loop( max_iterations=3 ) ``` Just a suggestion to show the difference and how to use override and partial together. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. Review Comment: ```suggestion results as a sequence (``LazyXComSequence``), the same as other downstream consumers of mapped tasks. ``` If I understood correctly. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. + +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``. + +Connect a loop to other tasks +------------------------------ + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + start() >> refinement >> finish() + +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. A task that must run in every iteration belongs inside the task group. Review Comment: ```suggestion for the loop. ``` I think that statement should be obvious? but feel free to ignore ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. + +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``. + +Connect a loop to other tasks +------------------------------ + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + start() >> refinement >> finish() + +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. A task that must run in every iteration belongs inside the task group. + +.. _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 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 + + 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. +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 default ``all_success`` rule prevents it from running, and +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. + +Skipping the body's terminal task also skips the gate under its default +``all_success`` trigger rule, so no next iteration is created. Downstream +tasks follow their own trigger rules; with the default, they are skipped too. Review Comment: ```suggestion tasks follow their own trigger rules. ``` I think adding this might be confusing, even if correct. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. + +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``. + +Connect a loop to other tasks +------------------------------ + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + start() >> refinement >> finish() + +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. A task that must run in every iteration belongs inside the task group. + +.. _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 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 + + 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. +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 default ``all_success`` rule prevents it from running, and +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. + +Skipping the body's terminal task also skips the gate under its default +``all_success`` trigger rule, so no next iteration is created. Downstream Review Comment: ```suggestion Skipping the body's terminal task also skips the gate under its ``all_success`` trigger rule, so no next iteration is created. Downstream ``` again this suggestion is if you cannot change the trigger rule of the gate task ########## airflow-core/docs/authoring-and-scheduling/loops-and-mapped-tasks.rst: ########## @@ -0,0 +1,61 @@ + .. 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-and-mapped-tasks: + +====================== +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, and sometimes neither is what you need. Review Comment: ```suggestion yet. Airflow has two features for this, as well as mechanisms like ``retries`` and dynamic Dag generation that fit adjacent use cases. ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks within an iteration can run in parallel. Review Comment: ```suggestion Task instances within an iteration can run in parallel. ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. Review Comment: ```suggestion failed. Meeting the condition on the last allowed iteration marks the gate task as success. ``` style suggestion to make the sentence easier to read ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``refine.loop(max_iterations=3)`` runs Review Comment: ```suggestion For a fixed-count loop, omit ``until``. The definition ``refine.loop(max_iterations=3)`` runs ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) Review Comment: ```suggestion @dag ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) Review Comment: ```suggestion @dag ``` I am on a personal mission to make agents understand that the catchup default is now False. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. Review Comment: ```suggestion 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``. ``` ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. + +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``. + +Connect a loop to other tasks +------------------------------ + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + start() >> refinement >> finish() + +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. A task that must run in every iteration belongs inside the task group. + +.. _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 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 + + 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. +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 default ``all_success`` rule prevents it from running, and +no next iteration is created. Review Comment: ```suggestion 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. ``` Is it possible to change the trigger rule of the gate? I assume not? If then I'd state that explicitely. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + Review Comment: ```suggestion ``` Moved this further down as a list of limitations. ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` Review Comment: ```suggestion Its return value (the XCom pushed under the `return_value` key) becomes ``loop.result`` for the gate and ``loop.previous`` ``` in case anyone uses a traditional operator here ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. + +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``. + Review Comment: ```suggestion 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. ``` Adding the two paragraphs from above. I did not understand what ", including authored Dag results, before it starts streaming status updates." means...? What is an authored Dag result? using `@result`? ########## airflow-core/docs/authoring-and-scheduling/loops.rst: ########## @@ -0,0 +1,339 @@ + .. 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 bound; the runtime condition +determines how many iterations are needed within that bound. + +Tasks 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 bound 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: + +* ``True`` means stop: the condition has been met. +* ``False`` means continue, provided another iteration is allowed. + +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. For a fixed-count loop, the gate continues until the +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(schedule=None, catchup=False, tags=["example"]) + 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() + +``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. + +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 succeeds. + +For a fixed-count loop, omit ``until``: ``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. + +.. code-block:: python + + @dag(schedule=None, catchup=False, tags=["example"]) + 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() + +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``. + +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 for the current iteration, used by the gate. + +Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped +task's position. Neither is a task try number. A task can have several tries +within one iteration. + +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. + +``include_prior_dates=True`` cannot select a loop iteration from another run. +Publish the result through a task outside the loop if later 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, including authored +Dag results, before it starts streaming status updates. + +The loop body must have exactly one terminal task definition before the gate. +Its return value becomes ``loop.result`` for the gate and ``loop.previous`` +for the next iteration. A mapped terminal task supplies its collection of +results using normal task-mapping semantics. + +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``. + +Connect a loop to other tasks +------------------------------ + +Use the object returned by ``.loop()`` in dependencies: + +.. code-block:: python + + start() >> refinement >> finish() + +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. A task that must run in every iteration belongs inside the task group. + +.. _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 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 + + 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. +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 default ``all_success`` rule prevents it from running, and +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. + +Skipping the body's terminal task also skips the gate under its default +``all_success`` trigger rule, so no next iteration is created. Downstream +tasks follow their own trigger rules; with the default, they are skipped too. +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. + +The gate does not independently wait for every task in the body. If the +terminal task uses ``one_success``, for example, it can finish while another +branch is still running. The gate then follows its normal trigger rule and +can start the next iteration while that branch continues in the earlier one. Review Comment: Question: I assume this means that loop iterations can run in parallel? What if iteration 2 takes a long time but ends up fulfilling the until condition while iteration 10 is already running, is iteration 10 stopped, does it complete, what if iteration 10 would fail the until condition, will there be an iteration 11? (I might be misunderstanding this 😅) -- 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]
