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

Airflow如何配置实现上游失败时任务标记为upstream_failed

Airflow DAG 分支触发规则实现方案

要实现你描述的状态判定规则,优先使用Airflow内置能力即可,不需要编写复杂的自定义逻辑。

最优方案(Airflow 2.0+ 适用)

直接将regular_task的trigger_rule参数设置为none_failed_min_one_success即可,该内置规则的判定逻辑完全匹配你的三个要求:

  • 任意上游sensor处于failed或upstream_failed状态时,regular_task会被直接标记为upstream_failed,无需等待其余上游任务执行完成
  • 所有上游sensor执行结束、无失败状态,但没有任何sensor执行成功(即两个sensor均为skipped状态)时,regular_task会被标记为skipped
  • 所有上游sensor无失败状态,且至少1个sensor执行成功时,regular_task会正常启动运行

该触发规则是Airflow 2.0版本正式发布的稳定能力,不需要新增任何额外任务,改完参数即可生效,覆盖绝大多数生产环境场景。

低版本兼容方案(Airflow 1.x 适用)

如果你使用的Airflow版本低于2.0,没有内置上述触发规则,可以通过新增轻量状态检查任务的方式实现相同效果,无版本兼容问题:

  1. 在两个sensor和regular_task之间新增一个PythonOperator作为状态检查节点,配置该节点的trigger_rule='all_done'
  2. 在检查节点的执行逻辑中拉取两个上游sensor的运行状态,根据状态对应抛出异常或正常返回
  3. 将regular_task的触发规则改回默认的all_success即可
    参考实现代码:
from airflow.operators.python import PythonOperator, get_current_context
from airflow.exceptions import AirflowSkipException, AirflowFailException
from airflow.utils.state import State

def check_upstream_sensor_status():
    context = get_current_context()
    task_instance = context["ti"]
    # 拉取两个上游sensor的运行状态
    sensor1_state = task_instance.get_state(
        task_id="branch_1_sensor",
        dag_id=task_instance.dag_id,
        execution_date=task_instance.execution_date
    )
    sensor2_state = task_instance.get_state(
        task_id="branch_2_sensor",
        dag_id=task_instance.dag_id,
        execution_date=task_instance.execution_date
    )

    # 任意sensor失败则抛出失败异常,下游regular_task会自动标记为upstream_failed
    if State.FAILED in (sensor1_state, sensor2_state) or State.UPSTREAM_FAILED in (sensor1_state, sensor2_state):
        raise AirflowFailException("上游分支sensor存在执行失败")
    # 两个sensor均被跳过则抛出跳过异常
    if sensor1_state == State.SKIPPED and sensor2_state == State.SKIPPED:
        raise AirflowSkipException("所有上游分支sensor均被跳过,无满足触发条件的分支")
    # 其余场景(至少一个sensor执行成功)正常返回,触发下游任务运行
    return

# 定义状态检查任务
status_check = PythonOperator(
    task_id="upstream_status_check",
    python_callable=check_upstream_sensor_status,
    trigger_rule="all_done"
)

# 调整DAG依赖关系
branch_1_sensor >> status_check
branch_2_sensor >> status_check
status_check >> regular_task

原有配置不生效的原因

  • 配置trigger_rule='one_success'时,规则逻辑只会统计上游成功任务的数量,只要上游存在失败任务且最终成功数不足1,任务会直接被标记为skipped,不会识别上游失败的场景
  • 配置trigger_rule='all_done'时,规则只要求所有上游执行结束就触发任务运行,不会区分上游是成功、失败还是跳过状态,因此会出现全跳过场景下任务仍启动的问题

内容的提问来源于stack exchange,提问作者megadarkfriend

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 09:09:20