You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.23 02:24:22