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

Airflow配置:TaskGroup串行执行及故障场景处理需求

Airflow TaskGroup串行执行及容错方案

核心实现思路

  • 直接以TaskGroup实例建立上下游依赖,确保前一个TG内所有任务执行完成后才启动下一个TG
  • 通过TriggerRule.ALL_DONE规则实现容错:即使当前TG执行失败,后续TG仍能正常触发
  • 依赖绑定到TG整体而非单个任务,避免清理失败TG时重复执行后续TG

完整代码示例

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.utils.task_group import TaskGroup
from airflow.utils.trigger_rule import TriggerRule
from datetime import datetime

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

with DAG('serial_tg_fault_tolerant',
         default_args=default_args,
         schedule_interval='@daily',
         catchup=False) as dag:

    start = DummyOperator(task_id='start')

    # 第一个TaskGroup:Workflow_FRA
    with TaskGroup('Workflow_FRA') as tg_fra:
        run_task_fra = DummyOperator(task_id='run_task_FRA')
        run_next_task_fra = DummyOperator(task_id='run_next_task_FRA')
        run_task_fra >> run_next_task_fra

    # 第二个TaskGroup:Workflow_BEL
    with TaskGroup('Workflow_BEL') as tg_bel:
        run_task_bel = DummyOperator(task_id='run_task_BEL')
        run_next_task_bel = DummyOperator(task_id='run_next_task_BEL')
        run_task_bel >> run_next_task_bel

    # 第三个TaskGroup:Workflow_DEU
    with TaskGroup('Workflow_DEU') as tg_deu:
        run_task_deu = DummyOperator(task_id='run_task_DEU')
        run_next_task_deu = DummyOperator(task_id='run_next_task_DEU')
        run_task_deu >> run_next_task_deu

    end = DummyOperator(task_id='end')

    # 构建串行依赖 + 容错规则
    start >> tg_fra
    # 过渡任务确保前TG无论成败都触发下一个TG
    tg_fra >> DummyOperator(task_id='after_fra', trigger_rule=TriggerRule.ALL_DONE) >> tg_bel
    tg_bel >> DummyOperator(task_id='after_bel', trigger_rule=TriggerRule.ALL_DONE) >> tg_deu
    tg_deu >> end

关键配置说明

  1. TaskGroup串行执行
    直接将前一个TaskGroup实例(如tg_fra)作为上游,Airflow会自动等待该TG内所有任务完成后,才触发下游的过渡任务,进而启动下一个TG。

  2. 容错处理
    两个TG之间的DummyOperator设置了trigger_rule=TriggerRule.ALL_DONE,这个规则的作用是:不管上游TG的任务是成功、失败还是跳过,该过渡任务都会执行,从而保证后续TG能正常启动,不会被失败的TG阻断。

  3. 避免清理时重复执行后续TG
    后续TG的依赖是绑定在已成功执行的过渡任务上(如after_fra),当你清理失败的Workflow_BEL时,只有该TG内的任务会重新运行,after_fra和Workflow_DEU不会被重复触发——因为过渡任务已经处于成功状态,不会被清理操作影响。

简化版配置(无需过渡任务)

如果想精简代码,可以直接为TaskGroup的下游依赖设置触发规则,不过这种方式需要确保TG的入口任务应用了规则:

# 替代依赖构建部分
start >> tg_fra
tg_fra >> tg_bel
tg_bel >> tg_deu
tg_deu >> end

# 为每个TG设置触发规则,确保上游TG无论状态如何都能触发自己
tg_bel.trigger_rule = TriggerRule.ALL_DONE
tg_deu.trigger_rule = TriggerRule.ALL_DONE

这种方式更简洁,但需要注意:Airflow会将TG的触发规则应用到其内部的第一个任务,从而实现相同的容错效果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 00:47:32