Airflow生产环境中如何添加任务定义但不纳入DAG执行?
解决Airflow中添加任务但不纳入DAG的问题
以下是几种可行的方案,满足你添加任务定义用于代码评审,但不将其纳入现有DAG调度的需求:
方法1:将新任务定义放在DAG上下文外部
Airflow仅会把with DAG(...) as dag:代码块内定义的任务关联到该DAG。把t4、t5放在这个上下文之外,它们就不会被纳入DAG,也不会被自动调度执行。
import os from datetime import datetime from airflow import DAG from airflow.operators.empty import EmptyOperator # 新任务定义在DAG上下文外,不会被Airflow调度 t4 = EmptyOperator(task_id="task_4") t5 = EmptyOperator(task_id="task_5") with DAG( "partial", description="A simple partial DAG", start_date=datetime(2023, 1, 1), ) as dag: t1 = EmptyOperator(task_id="task_1") t2 = EmptyOperator(task_id="task_2") t3 = EmptyOperator(task_id="task_3") t1 >> t2 >> t3
后续需要集成这些任务时,只需将它们移到DAG上下文内,并添加对应的依赖关系即可。
方法2:使用TaskGroup并设置add_to_dag=False
如果希望新任务保持在DAG上下文内,但不被注册到DAG中,可以借助TaskGroup的add_to_dag参数,将其设为False。
import os from datetime import datetime from airflow import DAG from airflow.operators.empty import EmptyOperator from airflow.utils.task_group import TaskGroup with DAG( "partial", description="A simple partial DAG", start_date=datetime(2023, 1, 1), ) as dag: t1 = EmptyOperator(task_id="task_1") t2 = EmptyOperator(task_id="task_2") t3 = EmptyOperator(task_id="task_3") # 创建不添加到DAG的任务组,内部任务不会被调度 with TaskGroup("pending_tasks", add_to_dag=False) as pending_tasks: t4 = EmptyOperator(task_id="task_4") t5 = EmptyOperator(task_id="task_5") t4 >> t5 t1 >> t2 >> t3
后续启用任务时,只需将add_to_dag=False改为True,再把任务组关联到现有DAG的依赖链中即可。
方法3:临时设置任务的trigger_rule为NONE(不推荐)
通过给任务设置trigger_rule=TriggerRule.NONE,可以让任务永远不会被触发执行。但这种方法只是阻止任务运行,任务仍会出现在Airflow的DAG视图中,可能造成混淆,因此仅作为备选方案。
import os from datetime import datetime from airflow import DAG from airflow.operators.empty import EmptyOperator from airflow.utils.trigger_rule import TriggerRule with DAG( "partial", description="A simple partial DAG", start_date=datetime(2023, 1, 1), ) as dag: t1 = EmptyOperator(task_id="task_1") t2 = EmptyOperator(task_id="task_2") t3 = EmptyOperator(task_id="task_3") # 临时设置触发规则为NONE,任务不会被执行 t4 = EmptyOperator(task_id="task_4", trigger_rule=TriggerRule.NONE) t5 = EmptyOperator(task_id="task_5", trigger_rule=TriggerRule.NONE) t1 >> t2 >> t3
内容的提问来源于stack exchange,提问作者ItayB
相关产品推荐
相关产品推荐

