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

Airflow 1.x如何配置父任务失败运行子任务、上游失败时跳过的规则

Airflow 1.x 区分failed与upstream_failed状态的触发规则实现方案

Airflow 1.x 确实没有内置可以区分failed和upstream_failed状态的Trigger Rule,两种实现方案如下:

方案1:自定义Trigger Rule(推荐)

可以直接修改Airflow核心源码新增自定义Trigger Rule,一次修改全局可用:

  • 找到Airflow安装路径下的airflow/utils/trigger_rule.py文件
  • 新增自定义规则常量,同时加入合法规则集合
  • 调整状态判断逻辑匹配需求:所有父任务状态为success/failed时执行子任务,存在任意父任务为upstream_failed/skipped时跳过子任务

修改示例代码:

from airflow.utils.state import State

# 新增自定义规则常量
ALL_SUCCESS_OR_FAILED = "all_success_or_failed"

# 将新规则加入全局合法Trigger Rule集合
TRIGGER_RULES = {
    ALL_SUCCESS,
    ALL_FAILED,
    ALL_DONE,
    ONE_SUCCESS,
    ONE_FAILED,
    NONE_FAILED,
    NONE_SKIPPED,
    ALL_SKIPPED,
    DUMMY,
    ALL_SUCCESS_OR_FAILED  # 新增自定义规则
}

# 新增对应判断逻辑(可直接整合到调度器的状态判断流程中)
def judge_custom_trigger(upstream_states: set) -> bool:
    skip_states = {State.UPSTREAM_FAILED, State.SKIPPED}
    allowed_run_states = {State.SUCCESS, State.FAILED}
    if any(s in skip_states for s in upstream_states):
        return False
    return all(s in allowed_run_states for s in upstream_states)

修改完成后重启Airflow的scheduler和webserver服务,即可在DAG的任务定义中直接指定trigger_rule="all_success_or_failed"使用。

方案2:无侵入轻量替代方案

如果没有权限修改Airflow核心源码,可以通过ShortCircuitOperator实现相同逻辑,无需改动全局配置:

  • 给子任务前置一个状态判断任务,设置该判断任务的trigger_rule="all_done",保证无论上游任务状态如何,判断逻辑都能执行
  • 在判断逻辑中获取所有上游任务的状态,按需求返回布尔值控制下游任务是否执行

DAG示例代码:

from airflow.models import DAG
from airflow.operators.python_operator import ShortCircuitOperator, PythonOperator
from airflow.utils.state import State
from datetime import datetime

def check_upstream_status(**context):
    ti = context["ti"]
    dag_run = ti.get_dagrun()
    # 获取所有上游任务的状态
    upstream_states = [
        dag_run.get_task_instance(task_id).state
        for task_id in ti.task.upstream_task_ids
    ]
    skip_states = {State.UPSTREAM_FAILED, State.SKIPPED}
    allowed_states = {State.SUCCESS, State.FAILED}
    if any(s in skip_states for s in upstream_states):
        return False
    return all(s in allowed_states for s in upstream_states)

# DAG定义
with DAG(
    dag_id="custom_trigger_demo",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    # 示例父任务
    parent_task = PythonOperator(
        task_id="parent_task",
        python_callable=lambda: 1 / 0  # 模拟任务运行失败场景
    )

    # 上游状态判断任务
    upstream_check = ShortCircuitOperator(
        task_id="upstream_check",
        python_callable=check_upstream_status,
        provide_context=True,
        trigger_rule="all_done"
    )

    # 实际要执行的子任务
    child_task = PythonOperator(
        task_id="child_task",
        python_callable=lambda: print("子任务正常执行")
    )

    parent_task >> upstream_check >> child_task

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 15:54:08