Airflow DAG任务分组为空问题:如何将任务正确加入分组?
Airflow任务分组为空的解决办法
我有一个Airflow DAG,依赖关系为start_task >> group1、start_task >> group2。任务存储在tasks字典中用于填充DAG的任务分组,但运行后分组为空,代码如下:
from airflow import DAG from datetime import datetime from airflow.operators.empty import EmptyOperator from python.etl.dag_param_generator import DAGParamGenerator from airflow.utils.task_group import TaskGroup dag_param_generator = DAGParamGenerator('TEST_DINAMIC_DAG') tasks = { "task_1": EmptyOperator(task_id="task_1"), "task_2": EmptyOperator(task_id="task_2"), "task_3": EmptyOperator(task_id="task_3"), "task_4": EmptyOperator(task_id="task_4") } with DAG( dag_id=dag_param_generator.get_dag_params('dag_id'), schedule=dag_param_generator.get_dag_params('schedule'), start_date=dag_param_generator.get_dag_params('start_date'), catchup=dag_param_generator.get_dag_params('catchup') ) as dag: start_task = EmptyOperator(task_id="start_task") with TaskGroup("group1") as group1: for key, value in tasks.items(): if key in ['task_1', 'task_2']: print(value) globals()[key] = value with TaskGroup("group2") as group2: for key, value in tasks.items(): if key in ['task_3', 'task_4']: print(value) globals()[key] = value start_task >> group1 start_task >> group2
问题原因
当前写法无法将任务关联到TaskGroup的核心原因:
- 提前在tasks字典中创建的
EmptyOperator实例,此时还不属于任何DAG或TaskGroup上下文。 - Airflow的TaskGroup依赖上下文管理绑定任务——只有在
with TaskGroup(...)代码块内创建的任务,才会自动归属到该分组。 - 使用
globals()[key] = value仅能将变量存入全局命名空间,无法修改任务的归属关系。
修复方案
方案1:在TaskGroup上下文内创建任务(推荐)
直接在TaskGroup的with块中实例化任务,任务会自动关联到当前分组:
from airflow import DAG from datetime import datetime from airflow.operators.empty import EmptyOperator from python.etl.dag_param_generator import DAGParamGenerator from airflow.utils.task_group import TaskGroup dag_param_generator = DAGParamGenerator('TEST_DINAMIC_DAG') with DAG( dag_id=dag_param_generator.get_dag_params('dag_id'), schedule=dag_param_generator.get_dag_params('schedule'), start_date=dag_param_generator.get_dag_params('start_date'), catchup=dag_param_generator.get_dag_params('catchup') ) as dag: start_task = EmptyOperator(task_id="start_task") with TaskGroup("group1") as group1: # 在group1上下文内创建任务,自动归属到group1 task_1 = EmptyOperator(task_id="task_1") task_2 = EmptyOperator(task_id="task_2") with TaskGroup("group2") as group2: # 在group2上下文内创建任务,自动归属到group2 task_3 = EmptyOperator(task_id="task_3") task_4 = EmptyOperator(task_id="task_4") start_task >> group1 start_task >> group2
方案2:手动绑定任务到TaskGroup(适用于需复用任务字典的场景)
如果必须复用提前定义的tasks字典,可以通过修改任务的task_group属性,将任务绑定到指定分组:
from airflow import DAG from datetime import datetime from airflow.operators.empty import EmptyOperator from python.etl.dag_param_generator import DAGParamGenerator from airflow.utils.task_group import TaskGroup dag_param_generator = DAGParamGenerator('TEST_DINAMIC_DAG') tasks = { "task_1": EmptyOperator(task_id="task_1"), "task_2": EmptyOperator(task_id="task_2"), "task_3": EmptyOperator(task_id="task_3"), "task_4": EmptyOperator(task_id="task_4") } with DAG( dag_id=dag_param_generator.get_dag_params('dag_id'), schedule=dag_param_generator.get_dag_params('schedule'), start_date=dag_param_generator.get_dag_params('start_date'), catchup=dag_param_generator.get_dag_params('catchup') ) as dag: start_task = EmptyOperator(task_id="start_task") with TaskGroup("group1") as group1: for key in ['task_1', 'task_2']: # 手动将任务绑定到group1 tasks[key].task_group = group1 with TaskGroup("group2") as group2: for key in ['task_3', 'task_4']: # 手动将任务绑定到group2 tasks[key].task_group = group2 start_task >> group1 start_task >> group2
内容的提问来源于stack exchange,提问作者Oleksandr Zakharchenko
相关产品推荐
相关产品推荐

