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

Airflow DAG触发规则配置:精准控制End任务执行状态

解决方案

要实现你需要的逻辑,核心是区分四个并行任务自身执行失败和它们的上游任务失败导致这四个任务未正常执行这两种情况。之前用all_done触发规则无法区分,因为它只关心直接上游是否完成,不管失败原因。可以通过以下方式实现:

思路

  1. 让四个并行任务的上游任务(前置任务)保持默认的all_success触发规则,确保前置任务失败时,四个并行任务会被标记为upstream_failed而非进入执行状态。
  2. 用PythonOperator替代DummyOperator作为End任务,在任务逻辑中检查四个并行任务的状态:
    • 如果四个任务的状态都是success或failed(说明前置任务成功,它们正常执行了),则End任务成功。
    • 如果任何一个任务的状态是upstream_failed(说明前置任务失败,导致该任务未执行),则抛出异常让End任务失败,进而整个DAG失败。

代码示例

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.operators.python import PythonOperator
from airflow.utils.state import State
from airflow.utils.trigger_rule import TriggerRule
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 1, 1),
}

with DAG('target_dag', default_args=default_args, schedule_interval=None) as dag:
    # 前置上游任务
    pre_upstream_task = DummyOperator(task_id='pre_upstream_task')
    
    # 四个并行执行的任务
    parallel_task_1 = DummyOperator(task_id='parallel_task_1')
    parallel_task_2 = DummyOperator(task_id='parallel_task_2')
    parallel_task_3 = DummyOperator(task_id='parallel_task_3')
    parallel_task_4 = DummyOperator(task_id='parallel_task_4')
    
    # 设置依赖:前置任务完成后执行四个并行任务
    pre_upstream_task >> [parallel_task_1, parallel_task_2, parallel_task_3, parallel_task_4]
    
    # 定义End任务的检查逻辑
    def validate_parallel_tasks(**context):
        dag_run = context['ti'].dag_run
        # 获取四个并行任务的实例
        task_instances = [
            dag_run.get_task_instance(task_id='parallel_task_1'),
            dag_run.get_task_instance(task_id='parallel_task_2'),
            dag_run.get_task_instance(task_id='parallel_task_3'),
            dag_run.get_task_instance(task_id='parallel_task_4'),
        ]
        
        # 检查每个任务的状态
        for ti in task_instances:
            if ti.state == State.UPSTREAM_FAILED:
                raise ValueError(f"Task {ti.task_id} failed due to upstream failure, DAG aborted")
    
    # 定义End任务
    end_task = PythonOperator(
        task_id='End',
        python_callable=validate_parallel_tasks,
        trigger_rule=TriggerRule.ALL_DONE,
        provide_context=True,
    )
    
    # 设置End任务的上游依赖
    [parallel_task_1, parallel_task_2, parallel_task_3, parallel_task_4] >> end_task

逻辑说明

  • 当pre_upstream_task成功时,四个并行任务会正常执行(无论成功或失败),End任务会在它们全部完成后触发,检查到没有upstream_failed状态的任务,因此成功结束。
  • 当pre_upstream_task失败时,四个并行任务会被标记为upstream_failed,End任务触发后会检测到这个状态,抛出异常导致自身失败,整个DAG也会标记为失败。

内容的提问来源于stack exchange,提问作者snir.isl

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 11:27:38