kaxil commented on code in PR #73746:
URL: https://github.com/apache/airflow/pull/73746#discussion_r4217947774
##########
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:
+ task_group_dict = self.task_group.get_task_group_dict()
+ # With only the root group, the projection is the task graph, which is
acyclic here.
+ if len(task_group_dict) == 1:
+ return
+ cycles = [
+ f"{', '.join(cycle[:-1])} and {cycle[-1]}"
+ for task_group in task_group_dict.values()
+ for cycle in
task_group._find_dependency_cycles(group_dict=task_group_dict)
+ ]
+ if not cycles:
+ return
+ # Airflow 3.5 raises AirflowDagCycleException here instead; tracked at
+ # https://github.com/apache/airflow/issues/73678
+ warnings.warn(
+ f"Dag '{self.dag_id}': {'; '.join(cycles)} depend on each other in
a cycle. Cyclic TaskGroup "
Review Comment:
This lists every member of every cycle, and it ends up in
`DagWarning.message`, a plain `Text` column (64 KB on MySQL;
`DagCode.source_code` uses `Text().with_variant(MEDIUMTEXT(), "mysql")` for
this reason). The setup/teardown shape from the new docs with about 2,000 tasks
of ~30-character ids between `create` and `delete` gives a 68 KB message, and
under MySQL strict mode that insert fails at the `session.flush()` closing
`update_dag_parsing_results_in_db` (outside the warnings `try`, same
transaction as the serialized Dag writes), so the file's parse results wouldn't
save for a Dag this PR says keeps working. Could each cycle be capped at the
first N members plus "and K more", with a test on a large cycle?
##########
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
Review Comment:
Could this move to the module-level imports? Importing
`airflow.sdk.exceptions` doesn't load this module, so there's no cycle to
avoid, and `check_sdk_imports_in_core.py` honours `# noqa: SDK001` per line
wherever the import sits.
--
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]