邮件通知要拆成两层实现:底层是 SMTP 发送通道,负责把消息可靠投出去;上层是通知决策层,决定「哪些事件该发、合并成几封、发给谁」。避免通知过载靠四个动作——按指纹去重、按维度聚合成摘要、给同类通知设重复间隔、按级别分流渠道。热点监控(AI 热点监控工具)里一次平台接口异常可能同时命中几十个关键词,缺了决策层,收件箱几分钟就被塞满。下面给出两层的参数取值与可运行代码。
一、发送通道:端口与加密方式怎么选
SMTP 通道的起点是端口选择。三种常见组合的差异直接影响连接成功率:
| 端口 | 加密方式 | 连接流程 | 适用场景 |
|---|---|---|---|
| 25 | 无加密或 STARTTLS | 明文握手后可选升级 | 内网中继,公网多被封禁 |
| 465 | 隐式 TLS(SMTPS) | 建连即进入 TLS | 国内邮箱服务商常用 |
| 587 | 显式 TLS(STARTTLS) | 先明文再升级加密 | 云服务商与海外服务 |
选定端口后,认证凭据要放环境变量,不写进代码。多数邮箱服务商要求用「授权码」而非登录密码,且对单账号日发送量、单位时间连接数都有配额,超限会直接返回临时错误码。
1.1 异步发送与指数退避重试
监控系统的通知不能阻塞主循环。用 aiosmtplib 配合 asyncio 把发信改成协程,再用信号量压住并发连接数,避免被服务端判定为异常流量:
import asyncio
import os
import random
from email.message import EmailMessage
import aiosmtplib
SMTP_HOST = os.getenv("SMTP_HOST")
SMTP_PORT = int(os.getenv("SMTP_PORT", "465"))
SMTP_USER = os.getenv("SMTP_USER")
SMTP_PASS = os.getenv("SMTP_PASS")
_sem = asyncio.Semaphore(3) # 并发连接上限
RETRY_LIMIT = 3
def build_mail(to_addr: str, subject: str, body: str) -> EmailMessage:
msg = EmailMessage()
msg["From"] = SMTP_USER
msg["To"] = to_addr
msg["Subject"] = subject
msg.set_content(body) # 纯文本正文,摘要 + 详情链接
return msg
async def send_mail(to_addr: str, subject: str, body: str) -> bool:
msg = build_mail(to_addr, subject, body)
for attempt in range(RETRY_LIMIT):
try:
async with _sem:
await aiosmtplib.send(
msg,
hostname=SMTP_HOST,
port=SMTP_PORT,
username=SMTP_USER,
password=SMTP_PASS,
use_tls=(SMTP_PORT == 465),
start_tls=(SMTP_PORT == 587),
timeout=15,
)
return True
except (aiosmtplib.SMTPException, asyncio.TimeoutError):
if attempt == RETRY_LIMIT - 1:
return False # 转入死信记录,等人工或补偿任务处理
delay = min(2 ** attempt, 30) * random.uniform(0.5, 1.5)
await asyncio.sleep(delay)
退避里加随机抖动,多个通知同时失败时不会在同一时刻集体重试。重试次数留 3 次就够,无上限重试只会把服务端配额烧光。
1.2 通道健康度也要有数据
发送通道本身会失效:授权码被重置、服务商调整策略、出口 IP 进入退信名单。给每次投递记录状态码与耗时,按小时统计成功率,成功率跌破阈值时改走备用通道。备用通道可以是第二个发信账号,也可以是服务商提供的 HTTP 接口,切换逻辑写在发送函数外层,业务代码无感。
二、决策层:把 N 条热点压成 1 封邮件
通知过载的根源是把原始事件不加处理地一对一转成邮件。一次上游接口抖动,几十个关键词同时命中「采集失败」,逐条投递就是几十封邮件,真正需要处理的那条根因反而被埋掉。成熟告警系统的处理流水线值得照搬:去重 → 分组 → 抑制 → 静默 → 发送,每一步都在压缩通知数量。
2.1 指纹去重
给每条通知算一个指纹,取「事件类型 + 平台 + 关键词 + 标的 ID」这类稳定字段做哈希,不要把时间戳、热度数值这种每轮都变的字段算进去,否则去重永远命中不了。指纹既用于判重,也用于记录「上次发送时间」,两个用途共享同一个键,逻辑不容易走偏。
2.2 三个时间参数决定通知密度
| 参数 | 管的是什么 | 建议取值 | 设错的后果 |
|---|---|---|---|
| 聚合等待 | 新分组的首条事件等多久再发 | 普通 30 秒,紧急 0 秒 | 过长延迟关键通知,过短聚合失效 |
| 分组间隔 | 同组有新事件加入后的发送间隔 | 5 分钟 | 过短造成轰炸,过长感知滞后 |
| 重复间隔 | 未恢复的问题隔多久提醒一次 | 1 至 4 小时 | 过短形成疲劳,过长被遗忘 |
2.3 聚合缓冲区实现
import hashlib
import time
from collections import defaultdict
GROUP_WAIT = 30 # 秒
REPEAT_INTERVAL = 3600 # 秒
_buffer = defaultdict(list) # group_key -> [event, ...]
_first_seen = {} # group_key -> 首条时间
_last_sent = {} # fingerprint -> 上次发送时间
def fingerprint(evt: dict) -> str:
raw = f"{evt['type']}|{evt['platform']}|{evt['keyword']}"
return hashlib.md5(raw.encode("utf-8")).hexdigest()
def group_key(evt: dict) -> str:
return f"{evt['platform']}:{evt['type']}"
def offer(evt: dict) -> None:
fp = fingerprint(evt)
now = time.time()
if now - _last_sent.get(fp, 0) < REPEAT_INTERVAL:
return # 重复间隔内,只更新不重发
key = group_key(evt)
_buffer[key].append(evt)
_first_seen.setdefault(key, now)
def flush_ready() -> list:
"""由轮询循环每秒调用,取出到期分组"""
now, batches = time.time(), []
for key in list(_buffer.keys()):
if now - _first_seen[key] < GROUP_WAIT:
continue
events = _buffer.pop(key)
_first_seen.pop(key, None)
for e in events:
_last_sent[fingerprint(e)] = now
batches.append((key, events))
return batches
flush_ready 返回的每个分组渲染成一封邮件:标题写「平台 + 事件类型 + 命中数量」,正文列前 10 条摘要,其余给一个后台链接。一次网络抖动引发的 50 条事件,就变成一封「B 站接口连续超时,影响 50 个关键词」。
三、分级路由与频率上限
同一套通知不该走同一个出口。按严重级别分流,再给每个出口设硬上限,才能兜住配置失误:
- 给事件打级别标签:致命(采集全链路中断)、警告(单平台失败率升高)、提示(新热点上榜);
- 致命级聚合等待设 0 秒,立即发邮件并同步推 WebSocket 前端提醒;
- 警告级按平台维度聚合,5 分钟一封;
- 提示级不单发,进入每日摘要,固定时间一封;
- 在发送函数外再包一层计数器,单收件人每小时封数触顶后只累计不投递,恢复后补一封汇总。
抑制规则同样有用:采集调度器整体宕机时,把下游「单平台无数据」的警告全部压掉,只发根因那一条。抑制要限定生效范围——规则里必须指明「在哪些标签相同时才生效」,否则 A 平台的故障会顺手屏蔽掉 B 平台的真实问题。
静默是另一个开关,用在计划内动作上:升级依赖库、切换数据库、调试新数据源期间,手工设一个带时间范围的静默窗口,到期自动解除。静默范围要写具体,按事件类型加平台组合匹配,别用一个宽松条件把整个通知链路关掉后忘记恢复。
恢复通知不能省。只发故障、不发恢复,收件人无法确认问题是否已解决,会反复登录系统人工核对,等于把通知的价值又还回去了。
四、三个容易踩的坑
重试不做幂等是高频事故。发送超时未必代表没发出去,服务端可能已投递成功,盲目重试就是重复发信。给每封邮件生成唯一业务 ID,写入发送流水表并建唯一索引,投递前先查流水,存在成功记录就跳过。
其次是连接开销。每封邮件新建一次连接,握手与认证的成本会占掉大半时间,批量场景要复用会话,一次登录连续发多封。
再有是正文体积。把完整热点内容和分析结果全塞进邮件,容易触发服务商的内容检查与体积限制。正文只放摘要和链接,详情留在系统里。
常见问题(FAQ)
Q1:聚合窗口设多久合适?
普通通知 30 秒起步,致命级设 0 秒即时发送,兼顾时效与合并效果。
Q2:同一个热点反复触发怎么办?
用指纹记录上次发送时间,重复间隔内只更新缓存,不再重复投递。
Q3:邮件发不出去如何兜底?
指数退避重试三次,仍失败就写入死信表并降级为站内提醒和日志告警。