ashb commented on code in PR #74274:
URL: https://github.com/apache/airflow/pull/74274#discussion_r4223876067


##########
airflow-core/docs/authoring-and-scheduling/loops.rst:
##########
@@ -0,0 +1,324 @@
+ .. Licensed to the Apache Software Foundation (ASF) under one
+    or more contributor license agreements.  See the NOTICE file
+    distributed with this work for additional information
+    regarding copyright ownership.  The ASF licenses this file
+    to you under the Apache License, Version 2.0 (the
+    "License"); you may not use this file except in compliance
+    with the License.  You may obtain a copy of the License at
+
+ ..   http://www.apache.org/licenses/LICENSE-2.0
+
+ .. Unless required by applicable law or agreed to in writing,
+    software distributed under the License is distributed on an
+    "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+    KIND, either express or implied.  See the License for the
+    specific language governing permissions and limitations
+    under the License.
+
+.. _loops:
+
+=====
+Loops
+=====
+
+A loop repeats a group of tasks a number of times decided at runtime,
+using the results of each iteration. Each iteration can use the previous
+iteration's result. Use a loop for work such as refining an answer, following
+a hierarchy, or processing successive batches.
+
+A conditional loop works like a ``do … while`` loop: run the body, then test
+whether to run it again. Airflow's ``until`` condition expresses when to stop,
+so it corresponds to ``do { body } while not until(result)``. The body runs at
+least once. ``max_iterations`` sets an upper limit; the runtime condition
+determines how many iterations are needed within that limit.
+
+Iterations run one after another; task instances within an iteration can run 
in parallel.
+
+If you are not sure whether you need a loop or mapped tasks, see
+:ref:`loops-and-mapped-tasks`.
+
+Create a loop
+=============
+
+Define the body with ``@task_group`` and call ``.loop()`` on the decorated
+function. Supply a positive integer ``max_iterations`` to limit the number of
+iterations. An optional ``until`` callable decides when to stop early.
+
+Airflow evaluates ``until`` in a gate task downstream of the body's terminal 
task.
+It must return a ``bool``:
+
+* ``True`` means stop: the condition has been met.
+* ``False`` means continue, provided another iteration is allowed.
+
+Any other return value, including the ``None`` that a forgotten ``return`` 
produces, will result in
+the gate task failing, and the loop not continuing.
+
+Airflow creates the next iteration when the gate says to continue.
+It does not create all ``max_iterations`` up front:
+if the loop stops after four iterations, you see those four, without a tail
+of skipped iterations. A fixed-count loop runs exactly ``max_iterations``
+iterations.
+
+This example improves an estimate of the square root of two until the error
+is small enough:
+
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+   :start-after: [START refine_estimate]
+   :end-before: [END refine_estimate]
+
+``evaluate`` returns the result for the iteration. The gate reads it through
+``loop.result``. If another iteration runs, ``improve`` reads that same result
+through ``loop.previous``. The gate appears as a task in the loop, named after
+the condition function: ``refine.accurate_enough`` in this example.
+
+Tasks inside a loop body receive ``loop`` as a context parameter, so a 
parameter with that name in a
+loop body must be keyword-only and cannot have a default other than ``None``.
+
+In a templated field, a Jinja ``{% for %}`` block defines its own ``loop`` 
that hides this one for
+the length of the block. Read the Airflow value before the block and use that 
inside it, for example
+``{% set iteration = loop.index %}``.
+
+When ``until`` has no usable name, the gate is named ``__loop_gate`` instead.
+That covers a lambda, a ``functools.partial`` and a callable object. A gate 
name
+that matches a task in the loop body gets a ``__1`` suffix, as with any other
+task ID.
+
+Stopping because ``until`` returned ``True`` means the loop converged.
+Reaching ``max_iterations`` while the condition remains ``False`` means the
+loop did not converge: the gate task in the final iteration is marked as
+failed. Meeting the condition on the last allowed iteration marks the gate 
task as success.
+
+For a fixed-count loop, omit ``until``. The definition 
``refine.loop(max_iterations=3)`` runs
+three iterations, carrying results between them. Reaching the cap completes
+a fixed-count loop successfully. Its gate is named ``__loop_gate`` within the 
group.
+
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+   :start-after: [START fixed_loop]
+   :end-before: [END fixed_loop]
+
+For a task-group function with arguments, supply them with ``.partial()`` 
before
+calling ``.loop()``. Use ``.override()`` to configure the group, for example to
+give another loop a different ``group_id``:
+
+.. exampleinclude:: /authoring-and-scheduling/examples/example_task_loops.py
+   :start-after: [START partial_override_loop]
+   :end-before: [END partial_override_loop]
+
+Read the loop context
+=====================
+
+Declare ``loop`` as a keyword-only parameter on a task function. Airflow
+supplies it at execution time; leave it out when calling the task in the Dag
+definition.
+
+.. list-table:: Loop context
+   :header-rows: 1
+   :widths: 25 75
+
+   * - Attribute
+     - Meaning
+   * - ``loop.index``
+     - The current iteration number, starting at zero.
+   * - ``loop.max_iterations``
+     - The configured limit on the number of iterations.
+   * - ``loop.previous``
+     - The previous iteration's result, or ``None`` in iteration 0.
+   * - ``loop.result``
+     - The terminal body task's result (the XCom pushed under the key 
``return_value``) for the current iteration, used by the gate.
+
+Use ``loop.index`` for the loop iteration and ``ti.map_index`` for a mapped
+task instance's position. Each task instance can have several tries, each with 
its own ``ti.try_number``.
+
+Pass data between iterations
+============================
+
+Between iterations, use ``loop.previous``. Check ``loop.index == 0`` to handle
+the first iteration. Do not test ``loop.previous is None``: it is also ``None``
+when the previous terminal task returned nothing, and a previous result of
+``0``, ``False``, or an empty collection can be valid data.
+
+The loop body must have exactly one terminal task definition before the gate.
+Its return value (the XCom pushed under the ``return_value`` key) becomes 
``loop.result`` for the gate and ``loop.previous``
+for the next iteration. A mapped terminal task supplies its collection of
+results as a sequence (``LazyXComSequence``), the same as other downstream 
consumers of mapped tasks.
+
+If the body has several branches, finish with a task that combines their
+outputs. For example, if two branches end in ``refine_left`` and
+``refine_right``, add this inside the task group:
+
+.. code-block:: python
+
+   @task
+   def combine(left, right):
+       return {"left": left, "right": right}
+
+
+   combine(refine_left(), refine_right())
+
+``combine`` is now the single terminal task. The gate can read each result
+through ``loop.result["left"]`` and ``loop.result["right"]``; tasks in the
+next iteration use the corresponding keys in ``loop.previous``.
+
+Limitations:
+
+* ``include_prior_dates=True`` cannot select a loop iteration from another Dag 
run. Push the result to XCom through a task outside the loop if later Dag runs 
need to retrieve it with an ordinary XCom pull.

Review Comment:
   Mmmm yes.
   
   I think the only thing to do here is to make reading an xcom for a loop 
value "from outside the loop" to return a LazyXcomSequence.
   
   (There are lots of things that make this tricky or hard to use for various 
future typologies, but I think this is about the only thing we can do now)



-- 
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]

Reply via email to