kaxil commented on code in PR #73746:
URL: https://github.com/apache/airflow/pull/73746#discussion_r4210405066


##########
airflow-core/src/airflow/dag_processing/dagbag.py:
##########
@@ -491,7 +496,23 @@ def bag_dag(self, dag: DAG):
         :raises: AirflowDagCycleException if a cycle is detected.
         :raises: AirflowDagDuplicatedIdException if this dag already exists in 
the bag.
         """
-        dag.check_cycle()
+        from airflow.sdk.exceptions import TaskGroupCycleDeprecationWarning  # 
noqa: SDK001
+
+        self.task_group_cycle_warnings.pop(dag.dag_id, None)
+        with warnings.catch_warnings(record=True) as captured_warnings:
+            # DeprecationWarning is ignored by default outside __main__, which 
would hide it here too.
+            warnings.simplefilter("always", TaskGroupCycleDeprecationWarning)
+            dag.check_cycle()
+        for captured in captured_warnings:
+            if issubclass(captured.category, TaskGroupCycleDeprecationWarning):
+                self.task_group_cycle_warnings[dag.dag_id] = 
str(captured.message)

Review Comment:
   This records the warning before the cluster policies and `_add_to_bag` run, 
and `dag_warnings` only checks that the id is in `self.dags`. In a folder-wide 
bag (`airflow dags reserialize` passes `dagbag.dag_warnings` to the DB), an 
acyclic Dag bagged first followed by a cyclic duplicate with the same id leaves 
the acyclic Dag with a cycle warning, and the reverse order pops the real one. 
Keeping the message in a local and writing the dict after `_add_to_bag` 
succeeds would avoid both.



##########
airflow-core/docs/core-concepts/dags.rst:
##########
@@ -624,6 +624,59 @@ If you want to see a more advanced use of TaskGroup, you 
can look at the ``examp
 
     When using the ``@task_group`` decorator, the decorated-function's 
docstring will be used as the TaskGroups tooltip in the UI except when a 
``tooltip`` value is explicitly supplied.
 
+.. _concepts:taskgroup-cycles:
+
+Cyclic TaskGroup dependencies
+^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
+
+.. deprecated:: 3.4.0
+    Dags with cyclic TaskGroup dependencies are planned to fail Dag parsing 
from Airflow 3.5.
+
+When each TaskGroup is treated as a single unit, a TaskGroup and its siblings 
must not depend on each other
+in a cycle. A dependency into or out of any task in a group counts as a 
dependency of the whole group. This
+can make a group both upstream and downstream of a sibling, although no 
task-level dependency forms a cycle:
+
+.. code-block:: python
+
+    with TaskGroup("group1"):
+        a1 = EmptyOperator(task_id="a1")
+        a2 = EmptyOperator(task_id="a2")
+
+    with TaskGroup("group2"):
+        b1 = EmptyOperator(task_id="b1")
+        b2 = EmptyOperator(task_id="b2")
+
+    a1 >> b1  # group2 depends on group1
+    b2 >> a2  # group1 depends on group2
+
+A path that leaves a TaskGroup and comes back into it also forms a cycle, even 
when the tasks inside the

Review Comment:
   This rule also flags a common setup/teardown shape: `create` and `delete` in 
`TaskGroup("cluster")`, `work` outside it, and `create >> work >> 
delete.as_teardown(setups=create)` gives cluster -> work -> cluster. It matches 
the rule from the lazy-consensus thread, but since the planned rejection would 
make it a parse error, could the docs call out the setup/teardown case 
explicitly (and a test pin it), so users with that layout know it's intended?



##########
airflow-core/newsfragments/73746.significant.rst:
##########
@@ -0,0 +1,20 @@
+Deprecate cyclic TaskGroup dependencies
+
+A Dag whose TaskGroups depend on each other in a cycle, when each TaskGroup is 
treated as a single unit,
+now issues a ``TaskGroupCycleDeprecationWarning`` at parse time and shows a 
Dag warning in the UI that
+names the TaskGroups and tasks involved. These Dags still parse and run, 
because their task dependencies
+are acyclic, but Airflow 3.5 is planned to reject them at parse time.
+
+A dependency into or out of any task in a TaskGroup counts as a dependency of 
the whole group, so a path
+that leaves a TaskGroup and comes back into it is a cycle. For example, ``a >> 
bridge >> b``, with ``a``
+and ``b`` in the same TaskGroup and ``bridge`` outside it, makes the group 
depend on ``bridge`` and
+``bridge`` depend on the group, even if ``a >> b`` is also set.
+
+**Migration:**
+
+- Move tasks between TaskGroups, or out of them, so that each group depends on 
the others in one direction
+  only. In the example above, move ``bridge`` into the group, or move ``b`` 
out of it.
+- To catch these Dags in CI, run your Dag tests with
+  ``-W error::airflow.sdk.exceptions.TaskGroupCycleDeprecationWarning``.

Review Comment:
   `check_cycle()` only runs from `bag_dag`, so `-W error` catches these only 
in tests that load the Dags through a `DagBag` and assert no import errors; a 
test that imports the Dag file directly still passes. dags.rst states that 
condition, could the newsfragment carry it too?



##########
task-sdk/src/airflow/sdk/definitions/dag.py:
##########
@@ -1150,6 +1154,30 @@ def check_cycle(self) -> None:
                 f"Cycle detected in Dag: {self.dag_id}. Faulty task: 
{faulty_task_id}"
             )
 
+        self._warn_task_group_cycles()
+
+    def _warn_task_group_cycles(self) -> None:

Review Comment:
   The check only lives in the Python SDK's `check_cycle()`. Go SDK Dags 
serialize task groups too (`go-sdk/airflow/serialize.go`), but 
`lang_sdk_processor.py` only validates the serialized Dag, so they get neither 
this warning nor, as it stands, the planned parse-time rejection. Is that 
intended? If not, running the check on `SerializedTaskGroup` in core would 
cover every SDK; if Lang SDKs are out of scope, worth saying so in the docs and 
in #73678.



##########
task-sdk/src/airflow/sdk/definitions/taskgroup.py:
##########
@@ -702,6 +710,80 @@ def _sort_via_pass_numbering(
         sorted_indices = sorted(range(n), key=lambda i: (pass_of[i], i))
         return [nodes[i] for i in sorted_indices]
 
+    def _find_dependency_cycles(self, *, group_dict: dict[str, TaskGroup]) -> 
list[list[str]]:
+        """
+        Find children that depend on each other in a cycle when each child 
TaskGroup is one unit.
+
+        An edge into any task of a child group counts as an edge into the 
group, so a path that
+        leaves a group and comes back into it is a cycle even though the 
task-level graph is
+        acyclic. This is stricter than the ordering ``topological_sort`` 
needs, which only looks
+        at edges into a group's roots.
+
+        :return: one list of child node ids per cycle, each in insertion order
+        """
+        nodes = list(self.children.values())
+        id_to_idx = {nid: i for i, nid in enumerate(self.children)}
+        projected = [
+            self._project_upstream_ids(i, self._get_unit_upstream_ids(child), 
id_to_idx, group_dict)
+            for i, child in enumerate(nodes)
+        ]
+        members: dict[int, list[str]] = {}
+        for i, component in 
enumerate(self._find_projection_components(projected)):
+            members.setdefault(component, []).append(nodes[i].node_id)
+        return [node_ids for node_ids in members.values() if len(node_ids) > 1]
+
+    @staticmethod
+    def _get_unit_upstream_ids(child: DAGNode) -> Collection[str]:
+        if not isinstance(child, TaskGroup):
+            return child._topological_upstream_ids
+        upstream_ids = set(child._topological_upstream_ids)
+        upstream_ids.update(edge_id for task in child for edge_id in 
task.upstream_task_ids)

Review Comment:
   This walk goes through `MappedTaskGroup.__iter__` (here, and in 
`get_roots()` via `_topological_upstream_ids` on the line above), which raises 
`ValueError("Task-generated mapping within a mapped task group is not allowed 
with trigger rule 'always'")` for any direct child with `trigger_rule=ALWAYS`, 
even when the group is expanded over a literal list. So `@task_group def tg(p): 
t1(p)` with `t1` on `ALWAYS` and `tg.expand(param=[1, 2, 3])` parses on main 
but becomes an import error with this PR, which breaks the "they keep working" 
promise for Dags that have no cycle at all.
   
   Building the set from `child.upstream_task_ids`, the non-None 
`child.upstream_group_ids`, and the `upstream_task_ids` of each task in 
`child.iter_tasks()` avoids the validating iterator (roots are a subset of all 
member tasks, so nothing is lost). Could you add that shape to a test that 
asserts acyclic Dags don't warn? There's no negative test on the SDK side yet.



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