如何在Airflow中实现多条件任务触发?解决异常触发问题
问题解决:Airflow DAG错误触发下游任务
问题根源
当Task A执行失败时,Task B会进入upstream_failed状态(因上游失败未实际执行),而你设置的all_failed触发规则会将upstream_failed视为满足条件,导致D2被误触发。
解决方案
将D1和D2的trigger_rule从all_failed修改为one_failed。one_failed规则仅会在**上游任务自身执行失败(状态为failed)**时触发,不会把upstream_failed状态纳入判断范围。
修改后的完整代码
def do_op1_work(**kwargs): x = kwargs['dag_run'].conf.get('x') log.info('x: ' + str(x)) if x == 0: raise ValueError('Manual Exception') def do_op2_work(**kwargs): y = kwargs['dag_run'].conf.get('y') log.info('y: ' + str(y)) if y == 0: raise ValueError('Manual Exception') with DAG(dag_id='fulfill_uv', schedule_interval=None, default_args=default_args, catchup=False) as dag: op1 = PythonOperator(task_id='A', python_callable=do_op1_work, provide_context=True) op2 = PythonOperator(task_id='B', python_callable=do_op2_work, provide_context=True) op3 = DummyOperator(task_id='C') # 替换trigger_rule为one_failed op4 = DummyOperator(task_id='D1', trigger_rule='one_failed') op5 = DummyOperator(task_id='D2', trigger_rule='one_failed') op6 = DummyOperator(task_id='E', trigger_rule='one_success') # 简化依赖写法,更直观 op1 >> op2 >> op3 op1 >> op4 op2 >> op5 [op3, op4, op5] >> op6
补充说明
one_failed触发规则逻辑:只要有至少一个上游任务自身执行失败(状态标记为failed),当前任务就会启动,不会识别因上游失败未执行的upstream_failed状态。- 若使用Airflow 2.x版本,推荐将
provide_context=True替换为op_kwargs={'dag_run': '{{ dag_run }}'},适配新版本的参数传递规范。
内容的提问来源于stack exchange,提问作者karthik_varma_k
相关产品推荐
相关产品推荐

