kaxil commented on code in PR #74274: URL: https://github.com/apache/airflow/pull/74274#discussion_r4196706322
########## airflow-core/newsfragments/task-loops.feature.rst: ########## @@ -0,0 +1 @@ +Repeat decorated task groups with runtime loop conditions and selective clearing using ``.loop()``. Review Comment: Is Loops targeting 3.4.0? `v3-4-test` is being fast-forwarded to main for the betas right now (backports to it are paused), so once this merges, the fragment and both new pages go into the next 3.4.0 beta. If the `.loop()` stack doesn't land before rc1, the 3.4.0 release notes and docs would describe an API that isn't in the package. Moving the fragment into the PR that makes `.loop()` importable, named with that PR's number, would avoid that. It would also clear the red `check-newsfragment-pr-number` job, and stop towncrier rendering the entry as "(task-loops)" with no PR link. ########## 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. Review Comment: What happens with `until=lambda *, loop: ...`? Naming the gate after the function gives `refine.<lambda>`, and `validate_key` rejects that. A `functools.partial` or a callable instance has no `__name__`, and a condition with the same name as a body task collides with it. Should the gate fall back to `__loop_gate` when there's no usable name, or should the page say `until` has to be a named function? ########## 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 Review Comment: What does `finish` actually depend on here? The loop has a variable number of gate instances, and gates that said "continue" presumably end in success too, so "other trigger rules behave normally" doesn't say which upstream instances `one_success` or `all_done` are evaluated against. Related, with the straggler-branch case further down: once the last gate stops, does `finish` wait for a branch still running in an earlier iteration, and does that branch failing affect `finish`? Small one: the snippet at line 200 uses `start()` and `finish()`, which aren't defined anywhere on the page, and the first example calls its task `finished`. Reusing `refinement >> finished()` would keep it consistent. ########## 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 Review Comment: If `loop.previous` resolves the way an `XComArg` does, `None` doesn't only mean "first iteration". A missing `return_value` XCom already resolves to `None`, so a later iteration also sees `None` when the terminal task returned nothing or was marked success by hand. In the square-root example that silently restarts from the initial estimate. Would recommending `loop.index == 0` here (and in the example's `improve`) be safer? ########## 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 Review Comment: Does the gate pick up `retries` from `default_args` like any other task in the group? With a common `retries=3, retry_delay=5min`, a gate that fails because the cap was reached would retry against the same unchanged result and take about 15 minutes longer to fail. Probably worth saying whether non-convergence fails without retrying, and whether an exception in `until` retries. ########## 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 Review Comment: Is `loop` settled as the context name? I can see two collisions. Jinja's `{% for %}` binds its own `loop`, so in a templated field `{{ loop.index }}` inside a for block is Jinja's 1-based counter and `loop.previous` doesn't exist. I checked with a context `loop` whose `index` is 0: inside the block it renders 1, 2. If `loop` joins `KNOWN_CONTEXT_KEYS`, any existing `@task` with a `loop` parameter defaulting to something other than `None` starts failing at parse with "Context key parameter loop can't have a default other than None". It would also help to say what `loop` is for a task outside a loop. ########## 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. Review Comment: What happens when `until` returns something other than a bool? A forgotten `return` gives `None`. If that is treated as falsy, the loop runs to `max_iterations` and then fails as "did not converge", which points the user at the wrong cause. Saying either "must return a bool, anything else fails the gate" or "evaluated for truthiness" would settle it. ########## 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. Review Comment: How do setup/teardown tasks in the body fit the "exactly one terminal task" rule? `TaskGroup.get_leaves()` skips teardowns today, so if the gate hangs off the group's leaf, the gate and the next iteration's setup can start while the previous iteration's teardown is still running (recreating a cluster with the same name, say), and a teardown failure wouldn't stop the loop. Either describe that or add setup/teardown in the body to the unsupported list. ########## airflow-core/docs/authoring-and-scheduling/dynamic-task-mapping.rst: ########## @@ -16,25 +16,29 @@ under the License. .. _dynamic-task-mapping: +.. _mapped-tasks: -==================== -Dynamic Task Mapping -==================== +============ +Mapped tasks +============ -Dynamic Task Mapping allows a way for a workflow to create a number of tasks at runtime based upon current data, rather than the Dag author having to know in advance how many tasks would be needed. +Task mapping, also called dynamic task mapping, allows a way for a workflow to create a number of tasks at runtime based upon current data, rather than the Dag author having to know in advance how many tasks would be needed. This is similar to defining your tasks in a for loop, but instead of having the DAG file fetch the data and do that itself, the scheduler can do this based on the output of an upstream task. -Unlike a Python for-loop executed at DAG parse time, dynamic task mapping defers task creation until runtime, allowing the scheduler to determine the exact number of task instances based on upstream task outputs. +Unlike a Python for-loop executed at DAG parse time, task mapping defers task creation until runtime, allowing the scheduler to determine the exact number of task instances based on upstream task outputs. Right before a mapped task is executed the scheduler will create *n* copies of the task, one for each input. It is also possible to have a task operate on the collected output of a mapped task, commonly known as map and reduce. +.. seealso:: + Not sure whether you need mapped tasks or a loop? See :ref:`loops-and-mapped-tasks`. Review Comment: Further down this page (line 395, outside the diff) still tells people to use "task mapping methods for loops" inside task group functions. Now that Loops is its own feature, that sentence sends readers the wrong way. Something like "`expand()` to iterate over values" plus a pointer to `loops-and-mapped-tasks` would fit. ########## 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. + +Manually marking a gate successful completes it without evaluating ``until`` +or creating another iteration. This is an explicit override of normal gate +execution. Any later iterations retained after a selective clear remain unchanged. + +Iterations and execution history +-------------------------------- + +The loop view groups tasks and mapped instances by iteration. The gate's state +and logs explain why the loop continued, stopped, or failed. Each task's tries +and logs remain accessible within its iteration. + +Clear tasks inside a loop +-------------------------- + +Use these controls to select how far a clear extends through the loop: + +* **Clear downstream**, selected by default, includes downstream tasks. + Clearing a task this way can also clear the gate in the same iteration. +* **Clear later loop iterations**, selected by default, clears later + iterations when the selection includes a gate task, whether selected + directly or through **Clear downstream**. The gate runs again and decides + whether the loop should continue from that point. + +Clearing later iterations does not require the loop to run the same number of +iterations again. The gate can stop earlier or continue further, within the +configured limit. + +Rerun part of an iteration +~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Suppose each iteration contains this sequence: + +.. code-block:: text + + prepare → process → consume → gate + +The loop has already run iterations 0 through 4. You clear ``process`` in +iteration 2 with both default options selected: + +* ``prepare`` in iteration 2 remains completed. +* ``process``, ``consume``, and the gate in iteration 2 run again. +* Later iterations 3 and 4 are cleared. Whether replacement iterations run Review Comment: What happens to the existing tries and logs of iterations 3 and 4 if the gate now stops at 2? The history section says each task's tries and logs remain accessible within its iteration, but here those iterations stop existing. Are they kept as history, or removed? ########## 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 Review Comment: "Authored Dag results" will be opaque to most readers. Linking the feature would help, something like "including tasks marked with `@result` (see :ref:`dag-result`)". The same goes for "the experimental DagRun wait API" on the line above, which `dag-result.rst` also documents. ########## 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. + +Manually marking a gate successful completes it without evaluating ``until`` +or creating another iteration. This is an explicit override of normal gate +execution. Any later iterations retained after a selective clear remain unchanged. + +Iterations and execution history +-------------------------------- + +The loop view groups tasks and mapped instances by iteration. The gate's state Review Comment: Which view is "the loop view"? Is it the Grid, a new tab, or the task group's detail panel? Naming it, and saying how to get there, would let readers find it. ########## 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]) Review Comment: The prose above says "the next iteration can operate on a different collection of items", but this example maps over the constant `[1, 2]` in every iteration, so it doesn't show that. Building the list from `loop.previous` would. The block also needs a lead-in sentence; the paragraph before it ends on the zero-length expansion case. Small one at line 249: `[1, 2]` gives map indexes 0 and 1, so "mapped instance 2" doesn't exist in this example. "Mapped instance 1" would line up. ########## 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 Review Comment: Nit: this page uses an underline-only title with `-` sections and `~` subsections. `loops-and-mapped-tasks.rst` and `dynamic-task-mapping.rst` both use an overlined `=` title with `=` sections and `-` subsections. Matching the sibling scheme keeps the three pages consistent. ########## 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. + +Manually marking a gate successful completes it without evaluating ``until`` +or creating another iteration. This is an explicit override of normal gate +execution. Any later iterations retained after a selective clear remain unchanged. + +Iterations and execution history +-------------------------------- + +The loop view groups tasks and mapped instances by iteration. The gate's state +and logs explain why the loop continued, stopped, or failed. Each task's tries +and logs remain accessible within its iteration. + +Clear tasks inside a loop +-------------------------- + +Use these controls to select how far a clear extends through the loop: + +* **Clear downstream**, selected by default, includes downstream tasks. Review Comment: The existing dialog labels this option **Downstream**, not "Clear downstream" (`ui/public/i18n/locales/en/dags.json:84`). The default comes from the user's saved clear options, so it's only selected by default until they change it. Does the "clear later loop iterations" choice also exist on the REST clear endpoint and `airflow tasks clear`, and what does each default to? A CLI default that differs from the UI's would surprise 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: + +* ``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 Review Comment: Once `.loop()` exists, could these examples move into `airflow-core/src/airflow/example_dags/` and come in through `exampleinclude`, the way `dynamic-task-mapping.rst` does (lines 39 and 226)? Then the example-Dag import tests cover them and the page can't drift from the API. Right now the three runnable blocks have no output and nothing that runs them, which is fine for a spec. The implementation PR is the natural place to move them. ########## 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 Review Comment: This repeats the non-convergence rule from lines 100 to 102. It also has a stray line break at "with / the default". Keeping one statement and linking to it would avoid the two drifting apart. ########## 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 Review Comment: What does the Dag run end up as when the loop ends this way? A skipped gate stops the loop, but under `all_success` the downstream `finished` is skipped too, so a loop that ran out of items looks the same as one whose body was skipped by a branch. Is "ran out of work" meant to be a normal way to finish a loop, and if so should the page show the trigger rule that makes downstream run? There's a similar end-state question at line 283: after a gate is marked successful by hand, does downstream treat the loop as converged? -- 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]
