kaxil commented on code in PR #74353:
URL: https://github.com/apache/airflow/pull/74353#discussion_r4198297382
##########
task-sdk/src/airflow/sdk/definitions/_internal/loop.py:
##########
@@ -75,10 +75,9 @@ def execute(self, context: Context) -> None:
def create_loop(
factory: _TaskGroupFactory,
/,
- *args: Any,
+ *,
Review Comment:
Dropping `*args`/`**kwargs` here breaks
`test_loop_new_member_joins_live_pass_after_reserialization` in
`airflow-core/tests/unit/models/test_dagrun.py`, which still calls
`create_loop(body, max_iterations=3, add_member=True)` at line 209. That now
raises `TypeError: create_loop() got an unexpected keyword argument
'add_member'` in the second `dag_maker` block, before `verify_integrity` is
reached. Could that call become
`body.partial(add_member=True).loop(max_iterations=3)`, the same way
`test_loop.py` was migrated?
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -296,6 +298,38 @@ def converged(loop):
assert outside.state == State.SUCCESS
[email protected]
+def clear_loop_example_runs():
+ yield
+ clear_db_runs()
+ clear_db_xcom()
Review Comment:
`dag.test()` on a `DagBag`-loaded Dag goes through `sync_bundles_to_db` and
`sync_bag_to_db`, so after these three cases `refine_estimate`,
`fixed_task_loop` and `mapped_task_loop` stay behind in `dag`, `dag_version`
and `serialized_dag`. These module-level tests are not covered by the
class-level autouse cleanup further down the file. Could the teardown also call
`clear_db_dags()` (and `clear_db_serialized_dags()`)?
##########
task-sdk/src/airflow/sdk/definitions/decorators/task_group.py:
##########
@@ -125,6 +129,20 @@ def override(self, **kwargs: Any) ->
_TaskGroupFactory[FParams, FReturn]:
# TODO: FIXME when mypy gets compatible with new attrs
return attr.evolve(self, tg_kwargs={**self.tg_kwargs, **kwargs}) #
type: ignore[arg-type]
+ def loop(self, *, max_iterations: int, until: Callable[..., bool] | None =
None) -> TaskGroup:
+ """
+ Repeat this group's tasks, creating each iteration when the previous
gate continues.
+
+ The body must have one terminal task definition, which may be mapped.
+ Use ``partial()`` to supply body arguments and ``override()`` to
configure the group.
Review Comment:
Small thing with this guidance:
`body.partial(value=2).override(group_id="refine").loop(...)` builds the loop
correctly but emits `Partial task group 'body' was never mapped!`. `override()`
uses `attr.evolve`, which resets `_task_group_created` on the copy, so the
intermediate partial factory warns from `__del__` when it is collected. The new
test uses `override().partial()`, which is why it does not show up. The same
thing already happens with `expand()`, but since `.loop()` documents combining
the two, maybe either have `partial()`/`override()` hand ownership to the new
factory, or say which order to use? The "never mapped" wording also reads oddly
for a loop.
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -296,6 +298,38 @@ def converged(loop):
assert outside.state == State.SUCCESS
[email protected]
+def clear_loop_example_runs():
+ yield
+ clear_db_runs()
+ clear_db_xcom()
+
+
[email protected]("clear_loop_example_runs")
[email protected](
+ ("dag_id", "expected_iterations"),
+ [("refine_estimate", 4), ("fixed_task_loop", 3), ("mapped_task_loop", 2)],
+)
+def test_documented_task_loop_examples_complete(session, dag_id,
expected_iterations):
+ bag = DagBag(
+ dag_folder=str(Path(__file__).parents[3] /
"src/airflow/example_dags/example_task_loops.py"),
Review Comment:
Nit: `tests_common.test_utils.paths.AIRFLOW_CORE_SOURCES_PATH / "airflow" /
"example_dags" / "example_task_loops.py"` would avoid the `parents[3]`
arithmetic and the new `Path` import (`test_dag_parsing.py` does it that way).
Also the `and ti.working_set` filter below never drops anything, since
`get_task_instances` is an ORM select and the `do_orm_execute` listener already
restricts it to current attempts.
--
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]