Airflow如何配置仅在整个DAG运行失败时发送告警通知
Airflow仅在整个DAG运行失败时推送通知的实现方案
你之前用的任务级email_on_failure配置、单任务绑定的MS Teams webhook通知,触发逻辑都是绑定在单个任务上的,只要任意一个任务执行失败就会触发推送,和「DAG整体最终运行失败才通知」的触发时机不匹配,以下两种方案都可以实现需求:
方案1:配置DAG级别失败回调,增加最终状态校验
不要把失败回调绑定到单个任务上,直接在DAG初始化参数中配置on_failure_callback,同时在回调函数内增加DAG运行状态判断,只有当DAG Run的状态被正式标记为failed时,才调用邮件、Teams webhook的推送接口。
如果DAG配置了任务重试、失败分支兜底逻辑,状态校验可以避免任务首次执行失败进入重试时误发通知,示例代码片段:from airflow.utils.state import State def dag_failure_notify(context): dag_run = context.get('dag_run') # 仅DAG整体状态为失败时触发通知 if dag_run and dag_run.state == State.FAILED: # 替换为实际的通知推送逻辑 push_notify(f"DAG {dag_run.dag_id} 执行失败,运行时间:{dag_run.execution_date}") # DAG初始化时全局绑定回调,不要给单个任务配置 with DAG( dag_id="your_biz_dag", on_failure_callback=dag_failure_notify, # 其余DAG基础配置(调度周期、重试规则等) ) as dag: # 原有业务任务逻辑注意:如果DAG开启了补数、多DAG Run并发运行,回调内必须从入参
context中取当前触发事件对应的DAG Run实例,不要直接查询元数据库取最新记录,避免通知和实际失败的实例不匹配。方案2:新增独立末尾通知任务,配置特殊触发规则(生产环境优先推荐,逻辑更可控)
在现有DAG所有业务任务的最后,新增一个专门负责发失败通知的任务,做两个核心配置:- 将DAG内所有业务任务设置为该通知任务的上游,保证所有业务任务执行完成(无论成功/失败/跳过)后,才会走到通知节点
- 给通知任务设置触发规则为
TriggerRule.ONE_FAILED,即仅当所有上游任务中至少一个执行失败、且所有上游任务都执行完毕时才运行当前通知任务;如果上游全部执行成功,通知任务会自动跳过,不会发送误报。
示例代码片段:
from airflow.operators.python import PythonOperator from airflow.utils.trigger_rule import TriggerRule # 此处为原有业务任务:task_1、task_2、task_3... def send_failure_alert(): # 替换为实际的通知推送逻辑(邮件、Teams webhook等) print("推送DAG整体运行失败告警") failure_alert_task = PythonOperator( task_id="send_dag_failure_alert", python_callable=send_failure_alert, trigger_rule=TriggerRule.ONE_FAILED, dag=dag ) # 所有业务任务都作为通知任务的上游 [task_1, task_2, task_3] >> failure_alert_task这种方案完全不受中间任务重试、单任务失败后被兜底分支处理的影响,只有等整个DAG的所有任务执行完成、确定整体运行失败时才会触发通知,稳定性最高。
内容的提问来源于stack exchange,提问作者Victor Lefebvre
相关产品推荐
相关产品推荐

