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

Airflow如何配置Trigger Rule仅在父任务为success/failed状态时运行

Airflow 自定义触发逻辑实现方案

你需要的触发逻辑(父任务仅为success/failed终态、无upstream_failed状态时才运行任务)无法通过原生触发规则直接实现,可根据使用场景选以下两种落地方式:

方案1:零侵入快速实现(推荐单任务/少量任务场景)

不需要修改Airflow任何核心逻辑,直接组合原生all_done触发规则+前置状态校验即可:

  • 先将目标任务的触发规则设置为all_done,保证所有上游任务跑完后才会进入当前任务
  • 在任务执行逻辑的最前端增加上游状态校验,不符合条件就直接抛出AirflowSkipException跳过当前任务

参考代码(PythonOperator场景):

from airflow.exceptions import AirflowSkipException
from airflow.models import TaskInstance

def your_task_logic(**context):
    # 前置校验上游状态
    dag_run = context["dag_run"]
    current_task = context["task"]
    for upstream_tid in current_task.upstream_task_ids:
        upstream_ti = dag_run.get_task_instance(upstream_tid)
        upstream_state = upstream_ti.state
        # 存在upstream_failed直接跳过
        if upstream_state == "upstream_failed":
            raise AirflowSkipException("检测到上游任务为upstream_failed状态,跳过执行")
        # 上游状态不是success/failed直接跳过
        if upstream_state not in ("success", "failed"):
            raise AirflowSkipException(f"上游任务{upstream_tid}状态为{upstream_state},不符合触发要求")
    
    # 校验通过后写你的业务逻辑
    print("所有上游状态符合要求,开始执行业务")

# 任务定义时指定trigger_rule为all_done
your_task = PythonOperator(
    task_id="your_task",
    python_callable=your_task_logic,
    trigger_rule="all_done",
    dag=dag
)

这个方案兼容性极强,所有Airflow 2.x版本都能直接用,逻辑调整灵活,不需要额外配置。

方案2:自定义全局触发规则(推荐多任务复用场景)

如果有大量任务需要用到这个触发逻辑,可以自定义全局TriggerRule,不用每个任务都重复写校验代码:

  1. 编写自定义规则的判断逻辑,不要直接修改Airflow原生源码,避免版本升级丢失修改
from typing import Iterable
from airflow.models.taskinstance import TaskInstance
from airflow.utils.context import Context

# 自定义规则唯一标识
TRIGGER_RULE_NO_UPSTREAM_FAILED = "no_upstream_failed"

def check_no_upstream_failed(upstream_tis: Iterable[TaskInstance], context: Context) -> bool:
    all_finished = True
    has_upstream_failed = False
    has_invalid_state = False
    valid_terminal_states = {"success", "failed"}

    for ti in upstream_tis:
        state = ti.state
        if not state:
            all_finished = False
            break
        if state == "upstream_failed":
            has_upstream_failed = True
            break
        if state not in valid_terminal_states:
            has_invalid_state = True
            break
    
    return all_finished and not has_upstream_failed and not has_invalid_state
  1. 在Airflow启动时完成规则注册,把注册代码放到Airflow自动加载的plugins目录下即可,比如$AIRFLOW_HOME/plugins/custom_trigger_rule.py:
from airflow.utils.trigger_rule import TriggerRule, TRIGGER_RULE_CHECKERS
# 注册规则常量
TriggerRule.NO_UPSTREAM_FAILED = TRIGGER_RULE_NO_UPSTREAM_FAILED
# 注册规则校验逻辑
TRIGGER_RULE_CHECKERS[TRIGGER_RULE_NO_UPSTREAM_FAILED] = check_no_upstream_failed
  1. 后续任务定义时直接使用自定义规则即可:
your_task = PythonOperator(
    task_id="your_task",
    python_callable=your_business_logic,
    trigger_rule=TriggerRule.NO_UPSTREAM_FAILED,
    dag=dag
)

注意:注册逻辑必须在DAG加载前执行,放在plugins目录下Airflow会在启动时自动加载,不需要额外配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 13:18:19