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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 08:07:44