如何实现仅当直接上游任务失败时触发Airflow任务?
解决Airflow中仅当直接上游任务失败时触发后续任务的问题
问题场景
现有Airflow DAG的任务依赖如下:
task1 >> task2 >> [task3, task4]
期望执行逻辑:
- 先执行task1,再执行task2
- task2执行成功则触发task3
- task2执行失败则触发task4
使用TriggerRule.all_failed配置task4时,出现不符合预期的情况:当task1失败导致task2被标记为upstream_failed时,task4却被执行了——这是all_failed的预期行为(该规则会将upstream_failed视为失败类状态),但我们需要实现仅当直接上游task2自身执行失败(状态为failed)时才触发task4的逻辑。
解决方案
方案1:自定义TriggerRule
通过自定义触发规则,精准判断直接上游任务的状态是否为failed,排除upstream_failed等其他状态。
from airflow.utils.trigger_rule import TriggerRule from airflow.utils.state import State from airflow.models.taskinstance import TaskInstance class DirectUpstreamFailed(TriggerRule): def is_triggered(self, ti: TaskInstance, state): # 获取直接上游任务的状态字典 upstream_states = ti.get_direct_upstream_states() # 仅当直接上游任务状态为FAILED时触发 return any(status == State.FAILED for status in upstream_states.values()) and not any(status == State.UPSTREAM_FAILED for status in upstream_states.values()) # 给task4配置自定义触发规则 task4.trigger_rule = DirectUpstreamFailed()
方案2:使用分支任务(BranchPythonOperator)
通过分支任务判断task2的实际状态,动态决定是否执行task3或task4,逻辑更直观灵活。
from airflow.operators.python import BranchPythonOperator from airflow.utils.state import State def judge_task2_status(**context): # 获取task2的任务实例 task2_ti = context["dag_run"].get_task_instance(task_id="task2") if task2_ti.state == State.FAILED: return "task4" elif task2_ti.state == State.SUCCESS: return "task3" else: # 其他状态(如upstream_failed)下不触发任何后续任务 return [] # 创建分支任务 branch_task = BranchPythonOperator( task_id="branch_task", python_callable=judge_task2_status, provide_context=True, dag=dag ) # 调整依赖关系 task1 >> task2 >> branch_task >> [task3, task4] # 确保task3和task4仅在分支任务指向时执行 task3.trigger_rule = "all_success" task4.trigger_rule = "all_success"
方案对比
- 自定义TriggerRule:代码简洁,适合单一依赖的场景,可直接复用触发规则。
- 分支任务:逻辑更清晰,支持复杂的多分支判断,扩展性更强。
内容的提问来源于stack exchange,提问作者Vito De Tullio
相关产品推荐
相关产品推荐

