ashb commented on code in PR #74274: URL: https://github.com/apache/airflow/pull/74274#discussion_r4196879919
########## 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: I think KNOWN_CONTEXT_KEYS is always going to run into possible collisions no mater what we choose? But no, it's not fixed -- 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]
