kaxil commented on code in PR #73724:
URL: https://github.com/apache/airflow/pull/73724#discussion_r4210407836
##########
airflow-core/tests/unit/utils/test_task_group.py:
##########
@@ -1351,6 +1351,88 @@ def spy(self, nodes, projected):
assert position[f"r{i}"] < position[f"r{i + 1}"]
+def _make_sibling_groups_cycle():
+ with DAG("sibling_groups_cycle", schedule=None, start_date=DEFAULT_DATE)
as dag:
+ start = EmptyOperator(task_id="start")
+ 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")
+ end = EmptyOperator(task_id="end")
+ start >> [a1, b2]
+ a1 >> b1
+ b2 >> a2
+ [a2, b1] >> end
+ return dag
+
+
+def _make_group_bridged_by_outside_task():
+ with DAG("group_bridged_by_outside_task", schedule=None,
start_date=DEFAULT_DATE) as dag:
+ with TaskGroup("group"):
+ a = EmptyOperator(task_id="a")
+ b = EmptyOperator(task_id="b")
+ bridge = EmptyOperator(task_id="bridge")
+ a >> bridge >> b
+ return dag
+
+
+def _make_three_group_ring():
+ with DAG("three_group_ring", schedule=None, start_date=DEFAULT_DATE) as
dag:
+ groups = {}
+ for group_id in ("g0", "g1", "g2"):
+ with TaskGroup(group_id):
+ groups[group_id] = (EmptyOperator(task_id="first"),
EmptyOperator(task_id="second"))
+ groups["g1"][0] >> groups["g0"][1]
+ groups["g2"][0] >> groups["g1"][1]
+ groups["g0"][0] >> groups["g2"][1]
+ return dag
+
+
[email protected](
+ ("make_dag", "expected"),
+ [
+ pytest.param(
+ _make_sibling_groups_cycle,
+ {
+ None: ["start", "group1", "group2", "end"],
+ "group1": ["group1.a1", "group1.a2"],
+ "group2": ["group2.b1", "group2.b2"],
+ },
+ id="sibling-groups",
+ ),
+ pytest.param(
+ _make_group_bridged_by_outside_task,
+ {None: ["bridge", "group"], "group": ["group.a", "group.b"]},
Review Comment:
The deserializer sorts children by label, so the expected order here, in
`three-group-ring`, and in the two API tests is already the input order. If
`_sort_cyclic_projection` just returned `list(nodes)`, all of those would still
pass and only `sibling-groups` would fail. This is also the only case that
reaches the fallback through `_sweep_projection`, so nothing checks that the
sweep path places the cycle correctly among its other siblings.
A shape where label order is wrong would pin it: `group(a, b)` plus
top-level `bridge`, `after`, `x0`, `x1`, `x2`, wired `a >> bridge >> b >>
after`. The input order is `[after, bridge, group, x0, x1, x2]`, this PR
returns `[bridge, group, x0, x1, x2, after]`, and the identity version returns
the input unchanged. The nested `outer.inner` / `outer.bridge` shape from the
screenshots would be worth a param too, since every new case has the cycle at
the root.
##########
airflow-core/src/airflow/serialization/definitions/taskgroup.py:
##########
@@ -236,10 +236,13 @@ def topological_sort(
"""
Sort children topologically — a task always comes after its upstream
dependencies.
- See ``TaskGroup.topological_sort`` in task-sdk for the algorithm.
Cycles are
- treated as corrupt input: ``DAG.check_cycle`` rejects cyclic Dags
before
- serialization, so a cycle reaching this code indicates malformed
serialized data,
- and we raise ``ValueError`` rather than silently looping forever.
+ See ``TaskGroup.topological_sort`` in task-sdk for the algorithm.
Unlike the task-sdk
Review Comment:
The description says the Task SDK sort keeps raising so it can detect these
Dags at parse time, but #73746 adds its own `_find_dependency_cycles` for that,
and nothing at parse time calls the SDK `TaskGroup.topological_sort`. Its only
in-tree caller is the deprecated `DAG.topological_sort()`, which keeps raising
`AirflowDagCycleException` for Dags that pass `check_cycle()` and that #73746
only warns about. Should the SDK side get the same fallback so the two sorts
agree, or is the raise there intentional?
--
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]