如何在Airflow DAG中实现task1或task2失败时触发task4执行
如何在Airflow DAG中实现task1或task2失败时触发task4执行
嘿,这个需求在Airflow里其实很好实现,核心就是利用Airflow的**触发规则(Trigger Rule)**来控制task4的执行条件。我给你拆解一下具体步骤,附上代码示例,一看就懂~
首先,你需要明确:正常流程是task1 >> task2 >> task3,而task4需要在task1或者task2任意一个失败时启动,不管另一个任务的状态(比如task1失败后task2会被跳过,这种情况也要触发task4)。
具体实现步骤
- 导入Airflow的
TriggerRule类,这个类定义了任务执行的触发规则; - 给task4设置上游依赖为
task1和task2(这样task4能监听这两个任务的状态); - 给task4指定触发规则为
TriggerRule.ONE_FAILED——这个规则的含义是:只要至少有一个上游任务失败,当前任务就会执行,不管其他上游任务是成功、跳过还是其他状态,完美匹配你的需求。
完整代码示例
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.trigger_rule import TriggerRule from datetime import datetime # 模拟任务逻辑的函数 def task1_func(): # 这里可以写你的实际业务逻辑,比如故意抛出异常测试失败场景 # raise Exception("Task1 failed!") print("Task1 completed successfully") def task2_func(): # 同样可以测试失败场景 # raise Exception("Task2 failed!") print("Task2 completed successfully") def task3_func(): print("Task3 completed successfully") def task4_cleanup_func(): print("Running cleanup task4 due to task1 or task2 failure") # 定义DAG with DAG( dag_id="failed_task_trigger_cleanup", start_date=datetime(2024, 1, 1), schedule_interval="@daily", catchup=False ) as dag: task1 = PythonOperator( task_id="task1", python_callable=task1_func ) task2 = PythonOperator( task_id="task2", python_callable=task2_func ) task3 = PythonOperator( task_id="task3", python_callable=task3_func ) task4 = PythonOperator( task_id="task4_cleanup", python_callable=task4_cleanup_func, # 设置触发规则:至少一个上游失败则执行 trigger_rule=TriggerRule.ONE_FAILED ) # 设置任务依赖 task1 >> task2 >> task3 # 让task4同时依赖task1和task2,这样能监听两者的状态 [task1, task2] >> task4
测试场景验证
- 当task1失败时:task2会被标记为
skipped(因为上游失败),此时task4的上游有一个failed(task1)和一个skipped(task2),符合ONE_FAILED规则,task4会执行; - 当task2失败时:task3会被标记为
skipped,task4的上游有一个success(task1)和一个failed(task2),同样符合规则,task4会执行; - 当task1和task2都成功时:task4的上游全是
success,不会触发执行,符合预期。
这样就能完美实现你想要的逻辑啦~
备注:内容来源于stack exchange,提问作者Ab_sin
相关产品推荐
相关产品推荐

