如何在Airflow任务内为Airflow UI添加动态告警?
Airflow DAG列表动态告警实现方案
核心思路
借助Airflow插件系统结合Flask钩子机制,实时扫描DAG状态,在UI中动态弹出告警,替代官方文档里静态构建的方式。
具体步骤
1. 编写自定义插件
在Airflow的plugins目录下新建dag_alerts_plugin.py,代码如下:
from airflow.plugins_manager import AirflowPlugin from airflow.utils.session import create_session from airflow.models.dag import DAG from flask import request from flask_appbuilder import BaseView, expose, AppBuilder # 检查DAG问题的核心逻辑 def get_dag_alerts(): alerts = [] with create_session() as session: # 排查无调度配置的DAG no_schedule_dags = session.query(DAG).filter(DAG.schedule_interval.is_(None)).all() if no_schedule_dags: dag_ids = ', '.join([dag.dag_id for dag in no_schedule_dags]) alerts.append(f"⚠️ {len(no_schedule_dags)}个DAG未配置调度:{dag_ids}") # 排查暂停状态的DAG paused_dags = session.query(DAG).filter(DAG.is_paused.is_(True)).all() if paused_dags: dag_ids = ', '.join([dag.dag_id for dag in paused_dags]) alerts.append(f"⚠️ {len(paused_dags)}个DAG处于暂停状态:{dag_ids}") return alerts # 注册告警视图(可选,用于单独查看告警) class DAGAlertView(BaseView): @expose('/') def show_alerts(self): alerts = get_dag_alerts() return self.render_template('dag_alerts.html', alerts=alerts) # 主插件类 class DAGAlertsPlugin(AirflowPlugin): name = "dag_alerts_plugin" appbuilder_views = [ { "name": "DAG告警中心", "category": "Admin", "view": DAGAlertView() } ] # 向DAG列表页面注入告警 def on_appbuilder_init(self, appbuilder: AppBuilder): @appbuilder.app.after_request def inject_alerts(response): # 只在DAG列表页触发告警 if '/dags' in request.path and response.status_code == 200: alerts = get_dag_alerts() for alert in alerts: appbuilder.add_message(alert, 'warning') return response
2. 准备告警模板
在plugins下新建templates目录,再创建dag_alerts.html:
{% extends "appbuilder/base.html" %} {% block content %} <div class="page-header"> <h1>DAG告警列表</h1> </div> {% if alerts %} <div class="alert alert-warning"> {% for alert in alerts %} <p>{{ alert }}</p> {% endfor %} </div> {% else %} <div class="alert alert-success"> <p>所有DAG状态正常</p> </div> {% endif %} {% endblock %}
3. 生效验证
重启Airflow Webserver后:
- 打开DAG列表页面,存在问题的DAG会触发顶部告警提示
- 也可通过Admin菜单里的「DAG告警中心」查看完整告警列表
扩展提示
- 可在
get_dag_alerts函数中添加更多检查规则,比如排查无所有者的DAG、最近失败次数超阈值的DAG等 - 使用
create_session查询数据库是Airflow推荐的安全方式,可避免并发冲突
内容的提问来源于stack exchange,提问作者Tevett Goad
相关产品推荐
相关产品推荐

