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
关键配置说明
TaskGroup串行执行
直接将前一个TaskGroup实例(如tg_fra)作为上游,Airflow会自动等待该TG内所有任务完成后,才触发下游的过渡任务,进而启动下一个TG。容错处理
两个TG之间的DummyOperator设置了trigger_rule=TriggerRule.ALL_DONE,这个规则的作用是:不管上游TG的任务是成功、失败还是跳过,该过渡任务都会执行,从而保证后续TG能正常启动,不会被失败的TG阻断。避免清理时重复执行后续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
相关产品推荐
相关产品推荐

