Superset 用久了你会发现,它就像一个平时不闹脾气的老实人,数据看板挂墙上、同事用得很顺,一切岁月静好。但一旦某天早上你打开电脑,发现看板转圈、图表加载失败,或者用户群里有人喊"这个报表怎么打不开了",那种手忙脚乱、只能靠重启解决的感觉,真的不太舒服。更麻烦的是,Superset 出问题往往不是一下子全挂,而是先出现一些细微的征兆,比如日志里多了一些奇怪的报错、某个接口响应慢了一拍、数据库连接池悄悄被打满。如果我们能在这些征兆出现的第一时间就感知到,并且自动把消息推到钉钉或者邮件里,那局面就完全不一样了。

这篇文章就围绕我自己在项目里实际做的一套监控方案来聊,讲清楚怎么通过抓日志关键字和采集运行指标,让 Superset 的"小脾气"被及时发现,并且自动报警。整个过程不依赖复杂的监控平台,只用简单的 Python 脚本就能搞定,哪怕是不太熟悉运维的开发者,也能照着做出来。

一、问题背景:Superset 也会“生病”

Superset 是 Apache 旗下的开源数据可视化平台,通过它我们能快速做图表、写 SQL、搭仪表盘。它功能强,界面也好看,但背后依赖的组件不少:Web 服务、Celery 异步任务、Meta 数据库、Redis 缓存等等。任何一个环节出问题,都有可能让用户看到报错或者数据加载不出来。

实际运行中,我遇到过几种典型情况:

  • 连接池爆掉:数据库连接没有及时释放,导致新的查询拿不到连接,报 TimeoutError
  • Celery worker 卡死:异步导出报表的任务积压,队列越来越长,但没有人干活。
  • 内存打架:多个并发查询把内存吃满,Superset 进程被系统杀掉。
  • 日志刷屏:某条 SQL 语句触发了一个诡异的错误,日志里反复抛出 Traceback。

这些问题都有一个共同点:它们不会一下子让整个系统崩溃,而是先表现为日志异常和指标异常。如果我们只等着用户投诉,那已经晚了。所以,我们要在日志和指标里装上一个"哨兵",让它在问题苗头出现时立刻发出信号。

二、监控方案的整体思路

我的思路其实很朴素,可以分为两条线:

一条线是日志关键字采集。Superset 会把运行过程中的错误、警告、请求超时等信息写到日志文件里。我们定时去扫日志,找出像 ERRORTracebackTimeoutErrorConnectTimeout 这样的关键词,统计它们出现的频率,再看最近 5 分钟有没有突然增多。如果增多,说明系统正在经历某种异常。

另一条线是运行指标采集。日志是"事后记录",指标是"实时状态"。我们通过读系统自带的指标,比如 CPU 使用率、内存占用率、磁盘剩余空间,还有应用层面的指标,比如 Celery 队列长度、活跃连接数,来判断 Superset 是不是快要撑不住了。一旦超过预设阈值,同样触发告警。

两条线互相补充,日志能告诉我们是"哪里错了",指标能告诉我们"是不是马上要出事了"。再搭配一个告警推送模块,把异常消息发送到钉钉群或者邮箱里,运营和开发就能第一时间收到提醒。

整个方案只需要一台能运行 Python 脚本的机器(可以是 Superset 所在宿主机,也可以是另一台跳板机),不需要额外部署 Prometheus、Grafana 这些重型工具,非常轻量。

接下来我直接上实操。为了让示例统一,下面所有代码都用 Python 来实现,用到的库也都是 Python 自带的,比如 ostimeresmtplib,再外加上一个 requests 用来调 Webhook。

三、日志关键字的采集与告警实操

3.1 日志从哪里来

首先得确认 Superset 的日志落在哪里。通常有以下几种情况:

  • 如果是用官方 Docker 镜像部署的,日志会输出到 stdout,然后用 docker logs 查看,但这样不易解析。更合适的做法是在启动 Superset 前配置日志文件路径。
  • 如果是用虚拟环境手动部署的,日志一般会写在 logs/superset.loglogs/superset_error.log 里。

为了让我们自己的监控脚本能读到,我建议在 Superset 的配置里显式指定日志文件。比如在 superset_config.py 中加入这样的配置(这是一个配置示例,我们监控脚本读取它生成的日志文件):

# 设置日志文件路径
LOG_LEVEL = "INFO"
LOG_FORMAT = "%(asctime)s - %(levelname)s - %(message)s"
LOG_FILE = "/data/logs/superset/superset.log"

这样就保证了日志会持续写入同一个文件,我们监控脚本读起来就省事了。

3.2 用 Python 写一个日志扫描器

接下来写一个扫描脚本。它的核心逻辑是:

  • 记住上一次读取到的文件位置。
  • 从那个位置继续读取新增加的行。
  • 用正则匹配合适的关键字。
  • 统计每个关键字出现的次数,以及最近的错误日志片段。

下面是一个完整的日志扫描器示例:

import time
import re
import os

# 日志文件路径
LOG_FILE = "/data/logs/superset/superset.log"
# 关键字正则列表,匹配到就算一条异常
PATTERNS = [
    r"ERROR",
    r"Traceback",
    r"TimeoutError",
    r"ConnectTimeout",
    r"Connection refused",
    r"OperationError"
]

# 存放上一次读取的位置
offset = 0

# 如果文件不存在就先创建,避免报错
if not os.path.exists(LOG_FILE):
    open(LOG_FILE, "w").close()

# 获取当前文件大小,如果比 offset 小,说明日志被轮转过,需要从头读
file_size = os.path.getsize(LOG_FILE)
if file_size < offset:
    offset = 0

def scan_log():
    """扫描日志文件,返回异常关键字统计结果"""
    global offset
    
    # 重新打开文件,按 seek 定位
    with open(LOG_FILE, "r", encoding="utf-8", errors="ignore") as f:
        if offset != 0:
            f.seek(offset)   # 跳到上次读到的位置
        
        # 读取新增加的行
        new_lines = f.readlines()
        offset = f.tell()     # 记录当前读到哪了
        
        # 统计每个关键字出现的次数
        result = {}
        for line in new_lines:
            for pattern in PATTERNS:
                if re.search(pattern, line):
                    # 用关键字本身作为 key,否则用固定名称
                    key = pattern.strip("r\"'")
                    result[key] = result.get(key, 0) + 1
                    
                    # 顺便打印一下错误行,方便我们肉眼观察
                    # 这里只做一个最简单的输出,实际可以写到另一个文件
                    # print(line.strip())
        return result

# 测试一次扫描,输出结果
if __name__ == "__main__":
    result = scan_log()
    if result:
        for k, v in result.items():
            print(f"{k} 出现了 {v} 次")
    else:
        print("没有发现异常关键字")

这段代码能跑出基础统计,但想要达到"自动发现"的效果,还得跟告警联动。

3.3 把扫描结果变成告警

我们现在要做的是:如果某一次扫描中,ERROR 出现次数超过 10 次,或者 Traceback 出现次数超过 3 次,就立刻发一条告警到钉钉。代码如下:

import requests
import json

# 钉钉机器人 Webhook 地址(需要自己在钉钉群里添加机器人获得)
DING_WEBHOOK = "https://oapi.dingtalk.com/robot/send?access_token=你的token"

def send_dingtalk(message):
    """发送告警消息到钉钉群"""
    headers = {"Content-Type": "application/json"}
    data = {
        "msgtype": "markdown",
        "markdown": {
            "title": "Superset 异常告警",
            "text": "### Superset 异常告警\n\n" + message
        }
    }
    # 发送 POST 请求
    resp = requests.post(DING_WEBHOOK, headers=headers, data=json.dumps(data))
    if resp.status_code == 200 and resp.json().get("errcode") == 0:
        print("钉钉消息发送成功")
    else:
        print("钉钉消息发送失败", resp.text)

# 在刚才 scan_log 的基础上,判断是否需要告警
def check_and_alert():
    """扫描日志并判断是否触发告警"""
    result = scan_log()
    if not result:
        return
    
    alert_text = ""
    # 阈值设置:ERROR 大于 10 就告警,Traceback 大于 3 就告警
    threshold_map = {
        "ERROR": 10,
        "Traceback": 3,
        "TimeoutError": 5,
        "Connection refused": 2
    }
    
    for keyword, count in result.items():
        threshold = threshold_map.get(keyword, 5)  # 默认阈值是 5
        if count >= threshold:
            alert_text += f"- 关键字 `{keyword}` 在最近一分钟内出现了 {count} 次,超过了阈值 {threshold}。\n"
    
    if alert_text:
        send_dingtalk(alert_text)
    else:
        print("异常数量未超过阈值,不发送告警")

这样就完成了一个日志关键字的自动发现和告警闭环。当然,真实环境里不会只跑一次,而是要让它持续运行。我们可以用一个简单的 while 循环,每隔 30 秒扫描一次:

if __name__ == "__main__":
    while True:
        check_and_alert()
        time.sleep(30)   # 每 30 秒扫一次

这种循环方式虽然土,但非常稳定。相比用 cron,它是常驻进程,能更及时地响应日志变化。

四、指标采集与异常自愈

4.1 采集哪些指标

日志是在事情发生之后才留下痕迹,但有些问题在日志还没出现时就已经埋下了隐患。比如内存慢慢上涨,虽然没有产生任何报错,但等涨到临界点,系统会直接 OOM,连写日志的机会都没有。所以我们必须盯指标。

需要采集的指标大致分两类:

  • 系统级指标:CPU 使用率、内存占用率、磁盘剩余空间。
  • 应用级指标:Superset 的活跃请求数、数据库连接数、Celery 任务队列长度。

系统级指标可以通过 Python 的 psutil 库来获取,这个库很常用,能拿到几乎所有的系统状态。应用级指标则需要依赖 Superset 的一些内部 API,或者通过查询它的元数据库来获取。

4.2 用 Python 抓取指标

我们先来写一个抓取系统级指标的示例,用 psutil 库:

import psutil

# 注意:需要先用 pip install psutil 安装这个库
# 安装命令:pip install psutil

def get_system_metrics():
    """获取系统级指标,返回一个字典"""
    metrics = {}
    
    # 获取 CPU 使用率(百分比),interval=1 表示采集 1 秒内的平均值
    cpu_percent = psutil.cpu_percent(interval=1)
    metrics["cpu_percent"] = cpu_percent
    
    # 获取内存信息,返回一个 namedtuple,包含 total、available、percent 等
    memory_info = psutil.virtual_memory()
    metrics["memory_percent"] = memory_info.percent
    
    # 获取磁盘使用率,路径根据实际情况填写,比如 /data
    disk_usage = psutil.disk_usage("/data")
    metrics["disk_percent"] = disk_usage.percent
    
    return metrics

# 打印一下结果
if __name__ == "__main__":
    m = get_system_metrics()
    print(f"CPU: {m['cpu_percent']}%")
    print(f"内存: {m['memory_percent']}%")
    print(f"磁盘: {m['disk_percent']}%")

除了系统指标,我们还要盯应用指标。以 Celery 为例,假设 Celery 的 worker 是跑在同一台机器上的,我们可以通过检查 Celery 的任务队列长度来判断有没有积压。但直接获取队列长度要装额外插件,这里提供一个偷懒的办法:通过 Superset 的日志来看,但由于我们已经在做指标采集,可以直接查 Superset 的元数据库,比如 query 表中有多少个查询处于 running 状态。使用 Python 自带的 sqlite3pymysql 都可以。

下面是一个示例,用 pymysql 查询 Superset 的数据库,得到当前正在执行的查询数量:

import pymysql

# 连接 Superset 的元数据库(注意这里是 MySQL,如果不是请替换)
conn = pymysql.connect(
    host="127.0.0.1",
    port=3306,
    user="superset",
    password="你的密码",
    database="superset",
    charset="utf8mb4"
)

def get_running_query_count():
    """获取当前状态为 running 的查询数量"""
    cursor = conn.cursor()
    # 在 superset 中,query 表记录了用户执行的 SQL,status 一般为 'running' 或 'failed'
    sql = "SELECT COUNT(*) FROM query WHERE status='running'"
    cursor.execute(sql)
    result = cursor.fetchone()[0]
    cursor.close()
    return result

# 输出当前正在执行的查询数量
if __name__ == "__main__":
    count = get_running_query_count()
    print(f"当前正在执行的查询数: {count}")

这个指标很有参考价值。如果长时间一直有几十个 running 查询,说明连接池可能被卡住了。

4.3 设置阈值和告警

拿到指标后,下一步就是设置阈值。阈值的设定并没有绝对标准,需要根据你机器的规格来调整。比如我的机器是 4 核 CPU、8G 内存,我一般这样设:

  • CPU 使用率超过 85% 持续 5 分钟,告警。
  • 内存使用率超过 90%,告警。
  • 磁盘使用率超过 80%,告警。
  • 正在执行的查询数超过 20,告警。

我们写一个综合检查函数,把指标判断和告警放在一起:

import time

# 阈值配置
TRIGGER_CPU = 85
TRIGGER_MEMORY = 90
TRIGGER_DISK = 80
TRIGGER_QUERY_COUNT = 20

def check_metrics_and_alert():
    """采集指标并判断是否触发告警"""
    metrics = get_system_metrics()
    content = ""
    
    # 逐项判断
    if metrics["cpu_percent"] > TRIGGER_CPU:
        content += f"CPU 占用率过高: {metrics['cpu_percent']}% (阈值 {TRIGGER_CPU}%)\n"
    if metrics["memory_percent"] > TRIGGER_MEMORY:
        content += f"内存占用率过高: {metrics['memory_percent']}% (阈值 {TRIGGER_MEMORY}%)\n"
    if metrics["disk_percent"] > TRIGGER_DISK:
        content += f"磁盘使用率过高: {metrics['disk_percent']}% (阈值 {TRIGGER_DISK}%)\n"
    
    # 查询运行中的查询数
    try:
        query_count = get_running_query_count()
        if query_count > TRIGGER_QUERY_COUNT:
            content += f"运行中的查询数过多: {query_count} (阈值 {TRIGGER_QUERY_COUNT})\n"
    except Exception as e:
        # 元数据库连接异常也要告警
        content += f"无法获取运行中的查询数: {str(e)}\n"
    
    # 如果内容不为空,说明有异常需要告警
    if content:
        send_dingtalk("### 指标告警\n\n" + content.replace("\n", "\n\n"))
    else:
        print("所有指标正常")

if __name__ == "__main__":
    while True:
        check_metrics_and_alert()
        time.sleep(60)   # 每分钟检查一次

到这里,指标采集和告警也完成了。但我们还有一点可以优化:当发现异常时,能不能做一些简单的自动恢复动作?比如发现 Celery worker 卡死了,就自动重启它。这种"自愈"操作要注意风险,不能乱来,但确实可以提高系统的稳定性。下面是一个简单的自愈示例:当运行中的查询数一直很多时,自动重启 Celery worker 进程。

import subprocess

def restart_celery_worker():
    """重启 Celery worker(这里假设进程名包含 superset_worker)"""
    print("尝试重启 Celery worker...")
    # 使用 subprocess 调用 shell 命令杀掉旧进程并启动新进程
    # 注意:实际生产环境建议使用 systemctl 或者 supervisor,这里只是演示
    subprocess.run(["pkill", "-f", "celery.*superset_worker"], capture_output=True)
    time.sleep(3)
    subprocess.run(["nohup", "celery", "-A", "superset.tasks.celery_app:app", "worker", "--pool=prefork", "--concurrency=5", "&"], capture_output=True)

不过,自动重启有风险,如果问题不是卡死,而是资源不足,重启反而会加剧。所以建议只在特定条件非常明显的时候使用,比如连续好几次检测到了同一个异常,才触发自愈。

五、告警推送的几种姿势

钉钉消息推送前面已经演示了,这里再补充一个邮件通知的方式,因为并不是所有团队都使用钉钉。邮件告警适合不用聊天工具的团队。

5.1 钉钉机器人

钉钉机器人的核心就是 Webhook,我们还可以让消息带上不同颜色的标签,比如说高优先级用红色,普通告警用黄色。钉钉的 markdown 消息也支持这些样式。前面已经写过了,这里不再重复。

5.2 邮件通知

用 Python 的 smtplib 发送邮件也很简单。我们需要准备一个发件邮箱的 SMTP 账号信息。下面是一个完整的邮件发送示例:

import smtplib
from email.mime.text import MIMEText
from email.header import Header

# 配置邮箱信息
SMTP_SERVER = "smtp.example.com"    # SMTP 服务器地址
SMTP_PORT = 465                     # SSL 端口
SENDER_EMAIL = "monitor@example.com"  # 发件人邮箱
SENDER_PASSWORD = "你的邮箱密码"      # 发件人邮箱密码或授权码
RECEIVER_EMAILS = ["ops@example.com", "dev@example.com"]  # 收件人列表

def send_email(subject, content):
    """发送纯文本邮件"""
    # 包装邮件内容,这里使用 plain 格式,但不支持 markdown,所以保持纯文本
    msg = MIMEText(content, "plain", "utf-8")
    # 设置邮件标题
    msg["Subject"] = Header(subject, "utf-8")
    msg["From"] = SENDER_EMAIL
    # 多个收件人以逗号分隔
    msg["To"] = ", ".join(RECEIVER_EMAILS)
    
    # 建立 SMTP 连接并登录
    try:
        # 用 SMTP_SSL 连接,如果端口是 587 则用 starttls
        server = smtplib.SMTP_SSL(SMTP_SERVER, SMTP_PORT)
        server.login(SENDER_EMAIL, SENDER_PASSWORD)
        # 发送邮件
        server.sendmail(SENDER_EMAIL, RECEIVER_EMAILS, msg.as_string())
        server.quit()
        print("邮件发送成功")
    except Exception as e:
        print("邮件发送失败:", e)

# 测试发一封邮件
if __name__ == "__main__":
    send_email("Superset 监控告警", "测试:CPU 使用率超过 85%")

邮件不受聊天软件限制,但不够即时。在实际使用中,可以把钉钉作为主要告警通道,邮件作为备份,两者都接入。

六、优缺点和注意事项

6.1 优点

这套方案的优点很明显:

  • 简单直接,不需要引入 Kafka、Prometheus、ELK 这样的大件,只要一个 Python 脚本就能跑起来。
  • 可定制性高,你想监控什么关键字、设什么阈值,改代码就行。
  • 没有额外依赖,只要能访问到日志和相应端口就能监控。
  • 告警渠道灵活,既可以发钉钉,也可以发邮件,甚至还可以发企业微信。

6.2 缺点

但也不能忽略它的短板:

  • 日志采集是基于轮询的,有延迟。如果日志突然爆炸,脚本可能跟不上。
  • 对指标阈值的设定需要经验,设得太高容易漏报,太低容易烦人。
  • 无法覆盖所有内部异常,比如页面 JS 报错、前端渲染失败这类问题,从后端日志和指标都察觉不到。
  • 脚本本身常驻运行,如果脚本自己崩了监控就没了。所以最好再配合一个 crontab 来定时检查监控进程是否存活,但这就涉及多技术栈了,我们的示例保持 Python,就不展开。

6.3 注意事项

在实际落地的过程中,有几个坑希望能帮你避开:

  • 日志轮转问题。superset.log 可能被 logrotate 轮转,原来的文件被改名成 superset.log.1,此时再读 superset.log 会从空文件开始,导致我们监控的 offset 错乱。一个简单的处理方式是在扫描前检查文件大小,如果发现文件大小比记录的位置小,就认为日志被轮转了,这时把 offset 重置为 0,从头计数。
  • 重复告警。同一个问题可能在连续扫描中反复触发,导致微信群被刷屏。解决方案是在发送告警后做一个"静默期",比如一个小时内同一关键字只告警一次。实现起来很简单,用一个字典记录每个关键字的最近告警时间即可。
  • 时区问题。日志中的时间如果和服务器本地时间不一致,可能导致我们基于时间窗口统计时出现偏差。最好统一使用 UTC 或者都改成服务器本地时间,并且注意夏令时。
  • 权限问题。如果 Superset 的日志目录对监控脚本运行的用户没有读权限,那就麻烦了。可以使用 chmod 赋予权限,或者把脚本的运行用户加到相应的组里。
  • 告警文案要友好。发送给钉钉群的消息最好包含时间、机器名、具体异常内容,不要把几十条堆上去。我们可以在发送前做一次聚合和摘要。

关于静默期去重,这里给一个简单的实现片段供参考:

# 记录上次告警时间的字典
last_alert_time = {}

def should_alert(keyword, cooldown_seconds=3600):
    """判断某个关键字是否还在静默期内"""
    import time
    now = time.time()
    last_time = last_alert_time.get(keyword, 0)
    if now - last_time >= cooldown_seconds:
        last_alert_time[keyword] = now   # 更新告警时间
        return True
    return False

你在调用告警时,可以先用 should_alert 过滤一次,这样就能有效避免消息轰炸。

七、总结

围绕 Superset 的异常监控,我们做了两件事:一是扫描日志里的关键错误,二是采集机器运行指标和数据库查询状态。两者结合,能在问题影响用户之前提前发现,并且通过钉钉和邮件自动通知给相关人员。整个过程使用 Python 实现,代码量不大,部署起来也不复杂,但能明显减少半夜被叫起来的次数。

在做这套监控的半年多时间里,它帮我成功抓到过连接池泄漏、Celery worker 阻塞、磁盘空间不足、SQL 查询死锁等好多问题。虽然不是高深的技术,但它做到了"细水长流地守护"。如果你也因为 Superset 偶尔闹脾气而头疼,不妨照着这个思路搭一套自己的监控小工具。从读日志,到抓指标,再到推告警,其中的每一样都可以根据你实际环境里的现象去调整。技术本身不难,难的是愿意花时间观察它、了解它、维护它。但这份付出是值得的,因为系统稳定了,我们才能睡个安稳觉。