buzz 这个项目,是我被消息通知碎片化逼出来的产物。手上一堆定时脚本、爬虫任务、服务器监控,最烦的从来不是脚本报错,而是脚本跑完你不知道结果——日志躺在那儿没人看,邮件偶尔进垃圾箱,钉钉群机器人改个关键词就哑火。干脆花了一个周末写了个极简的统一消息推送服务,名字就叫 buzz。核心功能一句话就能说清:对外只提供一个 HTTP API,往里扔一条消息,它自动帮你发到邮箱、钉钉群、飞书群或者任意自定义 Webhook,带重试、去重和通知级别控制。这篇文章就把这个项目从设计到部署完整拆开讲,适合想自己搭一套可靠通知系统的开发者、运维和个人项目爱好者参考。
1. 项目定位:buzz 要解决的核心问题
1.1 消息通知的碎片化痛点
先说清楚我为什么要自建。那段时间我大概有几类消息源:服务器磁盘告警要发邮件,爬虫跑完要在群里喊一声,数据批处理失败想直接弹到手机,还有些边缘设备的状态变化想推给家里人。结果每个渠道都有自己的接入方式,邮件有 SMTP 配置,钉钉要加自定义机器人,飞书要搞签名校验,iOS 想弹通知还得走第三方推送。
这些渠道单独用都不难,难的是它们互相之间没有任何统一入口。今天给脚本 A 加了钉钉通知,明天脚本 B 想要同样的通知又得复制一遍配置。更麻烦的是失败处理——邮件发送失败没有任何告警,钉钉机器人被限流之后消息直接丢了,你根本不知道通知本身是不是成功送达。
buzz 的出发点就是把这些渠道全部收拢到一个 API 后面。对于上游任务来说,它只需要知道一个 URL 和一个 Token,发一条 HTTP 请求就算完事;至于走哪个渠道、要不要重试、怎么去重,都交给 buzz 内部处理。这个思路有点像把分散的插线板全部接到一个 PDU 上,后面接什么设备,前面的人不用关心。
1.2 为什么还要自己实现一次
提到消息通知服务,现成的方案并不少。Pushover 是国外老牌,体验很好但要付费且服务器在境外;Bark 对 iOS 用户非常友好,可它默认服务器是公共实例;Server酱是国内开发者常用方案,免费版有限额;直接钉钉群机器人最简单,但功能太薄。
我做了个对比表格,方便你判断什么情况下值得自建:
| 方案 | 部署成本 | 稳定性 | 扩展性 | 适合自己的场景 |
|---|---|---|---|---|
| 邮件直发 | 低 | 中等 | 低 | 只做备份,不追求实时 |
| 钉钉/飞书机器人 | 极低 | 中等 | 低 | 团队内部快速接入 |
| Bark 公共实例 | 极低 | 中等 | 低 | iOS 个人推送 |
| Server酱 | 低 | 中等 | 中 | 个人开发者 |
| 自建 buzz | 中 | 高 | 高 | 多渠道统一、可控性强 |
自建的核心理由是可控。你不需要依赖某个第三方实例的额度,也不用担心别人服务器宕机导致你的告警石沉大海。buzz 部署到自己的服务器或者 NAS 上,所有配置都在自己手里,消息记录也保存在本地数据库,隐私这块也安心不少。
1.3 设计目标与不做的事
做任何项目都得先画边界。buzz 我明确要做的有这几件:提供一个干净的 HTTP API;支持多渠道并发分发;消息失败自动重试;按级别和标签做基本的去重;部署尽量简单,最好一条 Docker 命令跑起来。
我也明确不做什么。不做多用户系统,不搞复杂的权限角色,不做前端管理界面,不做消息数据的永久存储。为什么?因为它就是一个家庭或个人级别的通知管道,不是企业级的消息中间件。加了用户系统意味着要处理注册、登录、会话,复杂度翻倍;加管理界面意味着要写前端,维护成本更高。这些对我实际使用场景没有实质帮助,砍掉反而让核心链路更稳。
这种取舍挺重要的。我见过不少个人项目写着写着就膨胀成平台,最后连作者自己都不愿意维护。buzz 的定位就是一个 24 小时默默干活的水管工,而不是一个需要操作界面的控制台。
2. 架构设计:一个接口收消息,多个渠道发通知
2.1 整体分层与请求流转
buzz 的架构我刻意保持了三层:接入层、处理层、渠道层。接入层就是 HTTP API,负责收消息和鉴权;处理层负责把消息放进队列、做去重、按策略路由;渠道层是各种适配器,真正把消息推出去。
请求流转大概是这样的:
- 客户端 POST 一条消息到
/v1/message,带上 Token。 - 接入层校验 Token、检查消息格式,顺便做一下最简单的限流。
- 消息落到处理层,先去 SQLite 里查一下幂等键是否最近处理过,有就直接返回成功。
- 通过去重后,消息投递到内存队列,主循环拿到消息后按配置好的渠道列表逐个调用适配器。
- 每个渠道返回成功或者失败,失败的消息按规则重试,重试超过次数就记录为失败并尝试降级渠道。
为什么要用三层而不是直接请求里同步发送?最开始我确实写过同步版本,一条消息进来,挨个调用邮件、钉钉、飞书,整个请求耗时好几秒,上游脚本经常等到超时。改成队列异步发送之后,API 基本 20 毫秒以内就能响应,发送慢的渠道完全不影响上游。
2.2 消息模型设计
消息模型看起来简单,但实际上决定了整个系统的灵活性。buzz 的消息结构我设计为五个字段:title、body、level、tag和可选的idempotency_key。
level有info、warning、error三档,它影响两个东西:一是发送渠道的选择,比如error级别的消息除了发到群里,还可以额外触发邮件和手机推送;二是渠道内部的处理,比如钉钉机器人对warning以上消息可以 @ 特定人。这个字段让同一个 API 能承载多种告警语义,而不需要为每种级别单独建入口。
tag是消息的标签,我用它做两件事。第一是去重,相同 tag 在固定时间窗口内的重复消息只发一次;第二是路由,可以在配置里指定某些 tag 只走邮件,某些 tag 只走钉钉。举个例子,爬虫状态消息我设 tag 为crawler,只发到群里;服务器磁盘告警 tag 为alert,群和邮件都发。
idempotency_key是客户端主动提供的幂等键,这个更可靠。HTTP 调用有时候会超时,客户端不确定请求到底成没成功,通常会重试,如果没有幂等键,同一条消息就发了两遍。客户端生成一个唯一 ID,buzz 存在 SQLite 里,重试的时候直接返回之前的结果。
2.3 渠道适配器与插件化注册
渠道层我实现了一个很简单的适配器模式。每个渠道就是一个类,实现统一的send(notification)方法,返回布尔值表示成功与否。然后有一个注册表,通过配置里的名称找到对应的实例。
接口大致是这样:
class Channel(ABC): @abstractmethod def send(self, msg: dict) -> bool: """发送消息,成功返回 True,失败返回 False""" raise NotImplementedError发邮件就写一个EmailChannel,调 SMTP;发钉钉就写一个DingTalkChannel,用requestsPOST 到机器人 Webhook;发飞书也类似。注册表就是一个字典,配置里写"channels": ["email", "dingtalk"],系统就自动实例化对应的通道。
之所以要插件化而不是写死,是因为我自己后来就遇到了新需求:有一次我想把孙女的幼儿园通知转发到群提醒,后来又想给智能家居的异常状态加一个自定义 Webhook。每加一个渠道只需要新增一个类文件,在配置文件里加一行,完全不用动主逻辑。这正是这个项目最值得长期投入的设计。如果一开始就把钉钉逻辑写死在主流程里,后面每加一个渠道都要战战兢兢地改核心代码,那才是噩梦。
3. 核心实现:FastAPI 骨架与关键逻辑
3.1 服务端基础骨架
我选了 Python 加 FastAPI,原因很简单:异步支持好,代码量小,写接口效率极高。服务端文件结构如下:
buzz/ ├── app.py # FastAPI 应用入口 ├── channels/ # 渠道适配器 │ ├── __init__.py │ ├── email.py │ ├── dingtalk.py │ └── feishu.py ├── queue.py # 内存队列与发送循环 ├── store.py # SQLite 存储与去重 └── config.py # 配置加载入口文件app.py的核心代码不长:
from fastapi import FastAPI, Header, HTTPException from pydantic import BaseModel import uuid, time from queue import push_message from store import check_idempotency, record_idempotency app = FastAPI() class Message(BaseModel): title: str = "" body: str = "" level: str = "info" tag: str = "default" idempotency_key: str = "" @app.post("/v1/message") async def send_message(msg: Message, x_token: str = Header(..., alias="X-Token")): if x_token != app.state.token: raise HTTPException(status_code=401, detail="invalid token") if msg.idempotency_key: if check_idempotency(msg.idempotency_key): return {"status": "ok", "duplicated": True} record_idempotency(msg.idempotency_key) # 进入异步处理,立即返回 await push_message(msg.dict()) return {"status": "ok", "queued": True}Pydantic 的BaseModel直接做了参数校验,少写一堆解析逻辑。X-Token从请求头里取,没有就 401,这里先不展开安全部分,后面有专门章节。push_message 内部把消息丢到 asyncio 队列,由后台任务真正消费发送,这样 API 响应就非常快。
3.2 幂等与去重
幂等这块是很多人会忽略的地方,但对消息系统来说是命根子。上游脚本经常用 requests 发通知,网络抖动导致请求超时,脚本一重试,消息就重复了。半夜收到三条一模一样的磁盘告警,那体验相当酸爽。
我在存储层放了两个逻辑:一个是显式的idempotency_key查重,另一个是基于tag的滑动窗口去重。前者适用于客户端自己会生成唯一 ID 的场景,后者适合客户端懒得管、只是想防止同一事件频繁刷屏的场景。
核心实现依赖 SQLite,建表如下:
CREATE TABLE idempotency ( key TEXT PRIMARY KEY, created_at INTEGER NOT NULL ); CREATE TABLE messages ( id INTEGER PRIMARY KEY AUTOINCREMENT, tag TEXT NOT NULL, level TEXT NOT NULL, title TEXT, body TEXT, created_at INTEGER NOT NULL );滑动窗口去重的逻辑是:收到一条消息时,如果该tag在当前时间往前推 30 秒内已经出现过,则丢弃。实现不复杂,一个 SQL 就能搞定:
SELECT COUNT(*) FROM messages WHERE tag = ? AND created_at > ?窗口长度我设为 30 秒,对多数告警场景足够。有人可能觉得用 SQLite 做去重有点重,但我亲测下来完全没问题,单机每秒几十条消息轻松应对,根本到不了性能瓶颈。
3.3 发送队列与重试机制
发送队列我一开始想用 Celery + Redis,后来觉得为了这么个小项目引入一堆外部依赖太蠢了。用量级根本达不到要分布式队列的程度,Python 的asyncio.Queue完全够了。
队列模块简化版:
import asyncio, json queue = asyncio.Queue(maxsize=2000) async def worker(): from channels import get_channels channels = get_channels() while True: msg = await queue.get() for ch in channels: for attempt in range(3): try: ok = await ch.send(msg) if ok: break else: await asyncio.sleep(2 ** attempt) except Exception: await asyncio.sleep(2 ** attempt) else: # 全部重试失败,标记一条日志并降级 await fallback_notify(msg) async def push_message(msg: dict): if queue.qsize() >= 2000: # 队列满,直接丢弃最旧消息 queue.get_nowait() await queue.put(msg)2 ** attempt是指数退避,第一次失败等 1 秒,第二次等 2 秒,第三次等 4 秒。重试之间这个等待很关键,很多渠道瞬时失败其实过几秒就能恢复,如果立刻重试大概率还是失败。队列上限设为 2000,超过之后就丢最旧的消息,原因很简单:消息通知系统里新的告警永远比旧的重要,宁可丢一条旧消息也不能阻塞新的。
3.4 客户端接入:curl、Python、Shell
服务端写完,客户端接入必须简单到令人发指,否则还是没人用。buzz 的使用方式基本就是发 HTTP 请求,所以 curl 就是最标准的客户端。
最简单的发送:
curl -X POST http://your-server:8000/v1/message \ -H "X-Token: your-token" \ -H "Content-Type: application/json" \ -d '{"title":"磁盘告警","body":"根分区剩余空间低于10%","level":"error","tag":"alert"}'Python 项目里我一般不直接写 requests 调用,而是封装一个十行以内的小函数:
import requests def buzz(title, body="", level="info", tag="default"): requests.post( "http://your-server:8000/v1/message", headers={"X-Token": "your-token"}, json={"title": title, "body": body, "level": level, "tag": tag}, timeout=2, )Shell 脚本里更简单,配合 curl 直接用。我有个专门跑备份的 shell 脚本,末尾就一行:
curl -s -m 5 -X POST $BUZZ_URL \ -H "X-Token: $BUZZ_TOKEN" \ -d "{\"title\":\"备份完成\",\"body\":\"$BACKUP_FILE\",\"level\":\"info\",\"tag\":\"backup\"}" \ >/dev/null || echo "buzz notify failed"-m 5设置 5 秒超时很重要,虽然 buzz 响应很快,但万一服务挂了,我们也不想让备份脚本卡在处理通知上。
4. 部署与运维实战
4.1 Docker 化部署
消息服务这种常年挂在后台的东西,不上容器化不舒服。Docker 化部署的核心收益是环境隔离和秒级迁移,我最终交付的形态就是一个镜像加一条docker run命令。
Dockerfile 很短:
FROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY . . EXPOSE 8000 CMD ["uvicorn", "app:app", "--host", "0.0.0.0", "--port", "8000"]docker-compose.yml 会稍微完整一点,把数据目录挂载出来:
services: buzz: image: buzz:latest container_name: buzz ports: - "8000:8000" environment: BUZZ_TOKEN: "please_change_me" BUZZ_DB_PATH: "/data/buzz.db" BUZZ_CHANNELS: "email,dingtalk" SMTP_HOST: "smtp.example.com" SMTP_PORT: "465" SMTP_USER: "notify@example.com" SMTP_PASS: "your-smtp-pass" DINGTALK_WEBHOOK: "https://oapi.dingtalk.com/robot/send?access_token=xxx" volumes: - ./data:/data restart: unless-stopped环境变量统一承载所有配置,这个设计很重要。用环境变量而不是配置文件,意味着换了机器直接把 compose 文件拷过去就能跑,密钥也不会混在代码仓库里。
4.2 访问安全:Token 鉴权与限流
buzz 暴露在公网上之后,安全问题必须认真对待。我做了三层防护:Token 鉴权、IP 白名单和简单限流。
Token 是必选的,没有它等于把消息接口裸奔在公网,任何人都可以往你的钉钉群和邮箱里灌垃圾。Token 我用的是secrets.token_hex(32)生成的 64 位字符串,不要自己手输密码,要用随机生成。
限流这块我实现得非常轻量,在 FastAPI 里加一个依赖函数,按 IP 加 Token 双重维度计数,一分钟超过 60 次直接 429。
from slowapi import Limiter from slowapi.util import get_remote_address limiter = Limiter(key_func=get_remote_address) @app.post("/v1/message") @limiter.limit("60/minute") async def send_message(request: Request, msg: Message, x_token: str = Header(...)): # ...有人可能会问,一个单用户的内部工具,限流是不是多此一举?我的看法是:消息服务一旦被滥用,影响的不是你的 API 资源,而是你绑定的渠道资源。比如钉钉机器人在短时间内收到几百条消息,会触发平台限流,后续真正的告警反而进不来。限流不是为了防攻击,是为了避免上游脚本出 bug 刷爆渠道。
4.3 HTTPS 与反向代理
直接让 Uvicorn 对外暴露 HTTP 端口其实不太合适,一个是缺少 TLS,另一个是 Uvicorn 本身的静态文件处理能力有限。我习惯用 Caddy 做反向代理,它会自动申请和更新 HTTPS 证书,配置也简单。
Caddyfile 大概是这样的:
buzz.example.com { reverse_proxy 127.0.0.1:8000 }Caddy 会自动通过 ACME 拿到证书,HTTP 请求会被 301 跳到 HTTPS。这个方案相对 Nginx 加 certbot 来说少了很多手工步骤,尤其适合个人项目。加 HTTPS 之后,客户端代码里的请求地址都改成https://开头,Token 在传输过程中就不会被明文抓到了。
如果只是在内网用、不暴露公网,可以不做 HTTPS,但内网也要注意局域网内抓包风险。反正现在 Let's Encrypt 免费,建议一步到位。
4.4 数据持久化与备份
SQLite 文件是 buzz 唯一需要持久化的东西,挂载目录之后基本就完了。我用的目录结构是/data/buzz.db,你备份的时候只需要把这个文件复制走。
我写了个简单的每日备份任务,凌晨用sqlite3在线备份接口创建一个一致性快照:
sqlite3 /data/buzz.db ".backup '/backup/buzz-$(date +%F).db'" find /backup -name "buzz-*.db" -mtime +7 -deletesqlite3的.backup命令可以在数据库运行的时候安全执行,不会产生锁冲突,这点比直接 copy 文件靠谱得多。消息数据我保留最近若干天就够,因为通知是实时接收的,历史消息只做调试用,没必要无限堆积。
5. 典型接入场景:把 buzz 用起来
5.1 服务器监控告警
最先接入 buzz 的就是服务器监控。我有一台低配云主机和一台家里的 NAS,cron 里挂着几个健康检查脚本。以前磁盘告警靠邮件,手机端经常忽略;现在统一走 buzz,结果好太多。
磁盘空间检查脚本核心逻辑:
#!/bin/bash THRESHOLD=80 USAGE=$(df -h / | awk 'NR==2 {print $5}' | tr -d '%') if [ "$USAGE" -gt "$THRESHOLD" ]; then curl -s -m 5 -X POST "$BUZZ_URL" \ -H "X-Token: $BUZZ_TOKEN" \ -d "{\"title\":\"磁盘空间告警\",\"body\":\"根分区使用率 ${USAGE}%\",\"level\":\"error\",\"tag\":\"server\"}" fi我把BUZZ_URL和BUZZ_TOKEN写到 cron 环境变量里,所有脚本统一引用,换 Token 的时候只改一处。这个习惯很重要,避免 Key 散落在各种脚本里。
5.2 定时任务完成通知
cron 任务的通知我用的模式是:任务执行成功发 info,失败发 error。比如每天晚上的数据库备份任务,直接把输出压缩后通过 buzz 发出来,成功了报告备份文件名,失败了报告错误日志。
/usr/local/bin/backup.sh >> /var/log/backup.log 2>&1 if [ $? -eq 0 ]; then curl ... -d "{\"title\":\"备份成功\",\"body\":\"$BACKUP_FILE\",\"level\":\"info\",\"tag\":\"backup\"}" else curl ... -d "{\"title\":\"备份失败\",\"body\":\"$(tail -5 /var/log/backup.log)\",\"level\":\"error\",\"tag\":\"backup\"}" fi这里有个小技巧:把失败的关键几行日志塞到消息 body 里,这样人不用 ssh 进服务器排查,直接在群里就能看到大致原因。消息不用太长,截取尾部几行就够,避免刷屏。
5.3 爬虫与数据处理任务完成提醒
爬虫任务是最适合接通知的场景之一,因为批量爬虫经常在深夜跑,跑完不会有人一直盯着终端。我在 Python 爬虫里加了一个简单的装饰器,任务结束自动通知:
import functools from buzz_client import buzz def notify(fn): @functools.wraps(fn) def wrapper(*args, **kwargs): try: result = fn(*args, **kwargs) buzz(title=f"{fn.__name__} 完成", body=f"耗时 {result.elapsed}s,新增 {result.count} 条数据", tag="crawler") except Exception as e: buzz(title=f"{fn.__name__} 失败", body=str(e), level="error", tag="crawler") raise return wrapper这样一个装饰器就能让所有爬虫方法自动获得完成和失败通知,业务逻辑完全不用改动。大数据处理任务也一样,凌晨批处理跑完,第二天早上群里已经躺着结果摘要,相当于给自己配了个夜间值班助理。
6. 常见问题与排查技巧实录
6.1 常见问题速查表
实际运行中我整理了一些高频问题,直接做成了速查表,遇到问题可以对着查:
| 现象 | 可能原因 | 排查与解决 |
|---|---|---|
| API 返回 401 | Token 错误或请求头没带上 | 检查X-Token,确认环境变量是否正确加载 |
| 请求成功但渠道无消息 | 渠道配置了但适配器报错 | 看容器日志,逐渠道测试send方法 |
| 邮件经常收不到 | 端口或 SMTP 认证问题 | 检查 25/465/587 端口,确认 SMTP 密码未过期 |
| 钉钉机器人消息被拦截 | 机器人安全设置里的关键词限制 | 在钉钉后台添加与消息内容匹配的关键词 |
| 消息重复收到两次 | 上游超时重试且未带幂等键 | 客户端生成idempotency_key,或调整去重窗口 |
| 高并发时消息丢失 | 内存队列溢出丢弃旧消息 | 增加队列上限,或扩展单机资源 |
| 发送延迟大 | 某个渠道 SMTP 握手超时 | 给 SMTP 连接设置超时,重试放到后台线程 |
| 重启后历史消息消失 | SQLite 文件未持久化 | 检查 Docker 卷挂载是否正确 |
6.2 我踩过的几个坑和对应解法
第一个坑是邮件被丢进垃圾箱。SMTP 发送成功不代表邮件到达收件箱,我一开始用自己 VPS 的 IP 直发邮件,结果 Gmail 和 QQ 邮箱全都拦截了。后来换成域名邮箱的 SMTP 服务,又加了 SPF 记录,情况才好转。如果你只是做告警,建议直接用主流邮箱服务商提供的 SMTP,别自己搭邮件服务器找罪受。
第二个坑是钉钉机器人的关键词限制。钉钉自定义机器人有一个安全设置,要求消息内容里包含关键词才发送,如果 buzz 发出去的消息正文没命中关键词,钉钉会直接屏蔽。我当时在机器人后台加了一个不太常用的关键词,结果测试消息发不出来,排查了好久。解法很朴素:消息标题里统一加一个固定前缀,比如【监控】,然后在钉钉后台把监控设为关键词,这样所有告警都能放行。
第三个坑是时区问题。SQLite 里存时间戳我用的是 Unix 时间戳整数,这本身没问题,问题出现在客户端生成idempotency_key时用了本地时间字符串,导致不同时区的机器对同一条消息生成了不同的 key。后来统一要求客户端用 UUID 作为 key,彻底绕开时间概念。
第四个坑比较隐蔽:asyncio 事件循环里不小心写了同步阻塞调用。我早期在邮件适配器里用了同步的 SMTP 库,发送一封邮件偶尔要 2 秒,这个期间整个事件循环被卡住,其他消息全部排队。后来把同步调用放到asyncio.to_thread里执行,或者直接用aiosmtplib,问题才解决。经验是:FastAPI 的异步接口内部一定不要直接调用阻塞型库,否则再好的异步设计也白搭。
7. 写在最后:几个使用心得
buzz 跑到现在三个多月了,说几个真实体会供参考。最实用的用法反而不是服务器告警,而是给所有脚本、定时任务、批处理统一加通知,本质是给系统配了一个持续睁眼的口哨,哪条链路由问题,它第一时间吹响。另一个心得是渠道不要贪多,实际高频使用的往往就是邮箱加一个群机器人,其他渠道只是兜底方案;渠道配得越多,维护成本和出问题的面反而越大。最后建议你在接口层面尽早把幂等和限流做掉,这是我项目里返工最少的部分,但也是后来最庆幸先做的部分——很多问题就是在这两层被你挡住,后面才睡得安稳。如果你的通知需求也像我一样碎片化,与其在多个平台之间来回接线,不如照着这个思路,花一天时间搓一个属于自己的 buzz。