做实时日志系统的时候,告警功能几乎是标配。
最开始我的实现很简单粗暴:日志进来,逐条匹配规则,满足条件就发告警。
结果线上跑了一段时间,有一次某个服务疯狂报错,告警系统在两分钟内发出了 200 多条通知。运维的手机被叮个不停,最后直接把通知关掉了——反而错过了真正需要关注的窗口期。
这就是告警风暴(Alert Storm),也是大部分自研告警引擎最容易踩的坑。
这篇文章讲三件事:
- 用 Redis ZSET 实现滑动窗口计数,取代"每条日志都触发告警"的简单计数
- 用 Redis 冷却锁防止同一规则在短时间内重复告警
- 告警触发后,异步调用 Gemini 大模型自动完成根因分析,生成 Markdown 诊断报告
一、为什么不能用简单计数?
假设你有一条告警规则:
在 60 秒内,日志里出现
connection refused超过 10 次,就触发告警。
最直觉的实现是用 Redis INCR 计一个固定的计数器,每 60 秒重置一次。
问题在哪?固定窗口有边界效应。
比如计数器在第 59 秒重置,而错误恰好集中在第 58 秒到第 61 秒之间发生了 20 次:
|--- 窗口1 ---|--- 窗口2 ---|
... 10次 | 10次 ...
窗口 1 里 10 次,窗口 2 里 10 次,两个窗口都没超过阈值,告警压根没触发——但实际上短短 3 秒内出现了 20 次错误,已经是严重事故了。
滑动窗口解决的就是这个问题:它不是以固定时间点作为窗口边界,而是以"当前时刻往前推 N 秒"作为窗口,窗口随时间滑动,没有边界效应。
二、用 Redis ZSET 实现滑动窗口
Redis ZSET(有序集合)是实现滑动窗口的最优解:
- Key:唯一标识这条告警规则的计数桶
- Member:每次命中时写入一条记录(用 UUID 保证唯一性)
- Score:写入时的时间戳(毫秒)
有了这个结构,"窗口内有多少命中"这个问题就变成了:
在 ZSET 里,Score 在
(now - windowMs, now]范围内的元素有多少个?
ZSET key: alert:rule:42:hits
│
├── member: uuid-aaa score: 1720000001000 (01:00:01)
├── member: uuid-bbb score: 1720000003500 (01:00:03)
├── member: uuid-ccc score: 1720000058000 (01:00:58) ← 过期,需清理
└── member: uuid-ddd score: 1720000061000 (01:01:01)
代码实现:
go
// SlideWindowCount 基于 Redis ZSET 实现滑动窗口计数
// ruleID: 规则的唯一 ID
// windowSec: 时间窗口大小(秒)
// 返回:当前窗口内的命中次数
func SlideWindowCount(ctx context.Context, ruleID int, windowSec int) (int64, error) {
key := fmt.Sprintf("alert:rule:%d:hits", ruleID)
now := time.Now()
nowMs := now.UnixNano() / int64(time.Millisecond) // 当前时间戳(毫秒)
windowMs := int64(windowSec) * 1000 // 窗口大小(毫秒)
// 用 Pipeline 把三个操作打包,减少网络往返,同时保证操作的整体顺序
pipe := rdb.Pipeline()
// 1. 写入本次命中(Score = 当前时间戳,Member = UUID 保证唯一)
pipe.ZAdd(ctx, key, redis.Z{
Score: float64(nowMs),
Member: uuid.New().String(),
})
// 2. 清理已经滑出窗口的历史数据(Score < cutoff 的全部删掉)
cutoff := float64(nowMs - windowMs)
pipe.ZRemRangeByScore(ctx, key, "-inf", fmt.Sprintf("%v", cutoff))
// 3. 统计当前窗口内的元素数量
cardCmd := pipe.ZCard(ctx, key)
// 执行 Pipeline
_, err := pipe.Exec(ctx)
if err != nil {
return 0, err
}
return cardCmd.Val(), nil
}
这里有几个细节值得拿出来说:
为什么 Member 用 UUID,不用时间戳?
如果用时间戳作 Member,在同一毫秒内有两条日志同时命中,第二条的 ZAdd 会覆盖第一条(因为 Member 相同),导致少计一次。UUID 保证每次命中都是独立的 Member,不会相互覆盖。
为什么用 Pipeline,不用事务(MULTI/EXEC)?
Pipeline 把命令批量发送到 Redis,减少网络 RTT,性能更好。而这三个操作(ZAdd、ZRemRange、ZCard)不需要强事务保证——即使中途有别的命令插进来,也只是多了几条记录,不影响窗口计数的最终正确性。
ZRemRangeByScore 的必要性
如果不清理历史数据,ZSET 会无限增长,长期运行下去必然 OOM。每次命中时顺手清理掉窗口外的数据,Key 的大小就永远不会超过"窗口内最大命中次数",内存可控。
三、防止告警风暴:冷却锁
滑动窗口解决了"何时触发告警"的问题,但还没解决"触发了之后疯狂重复告警"的问题。
想象一下:一个 60 秒窗口、阈值 10 次的规则被触发了。接下来的每一条新日志进来,只要窗口内依然有 10 条以上,就会再次触发——也就是说,只要故障持续,告警会无限发。
解决方案很简单:告警触发后,用 Redis 设一把带过期时间的"冷却锁"。在锁的有效期内,这条规则不再重复发告警。
go
// IsAlertCooldown 检查这条规则当前是否在冷却期
func IsAlertCooldown(ctx context.Context, ruleID int) bool {
key := fmt.Sprintf("alert:rule:%d:cooldown", ruleID)
exists, err := rdb.Exists(ctx, key).Result()
if err != nil {
return false // Redis 故障时,降级放行(宁可多发一条,不能漏告警)
}
return exists > 0
}
// SetAlertCooldown 设置冷却锁,有效期内该规则不再重复告警
func SetAlertCooldown(ctx context.Context, ruleID int, duration time.Duration) error {
key := fmt.Sprintf("alert:rule:%d:cooldown", ruleID)
return rdb.Set(ctx, key, "locked", duration).Err()
}
注意 IsAlertCooldown 里的降级处理:Redis 出故障时返回 false(放行),而不是 true(拦截)。
这是一个关键的设计决策:冷却锁是为了减少噪音,但如果因为 Redis 故障就完全屏蔽告警,万一这时候线上真的出了严重问题,你反而什么都收不到。宁可多发几条重复告警,不能因为降级把真正的告警也干掉。
四、规则缓存:不能在日志热路径上查数据库
告警规则存在数据库里,但告警引擎是在日志处理的热路径上异步执行的——每一条进来的日志都可能触发 checkAlert,如果每次都去数据库查规则,QPS 高的时候数据库直接被打挂。
解决方案是在内存里维护一个本地规则缓存:
go
// CachedRule 包含一条告警规则,以及它预编译好的正则表达式
type CachedRule struct {
ID int
Name string
Pattern string // 原始正则字符串
Regexp *regexp.Regexp // ← 预编译好的正则句柄,直接复用
Source string // 来源过滤,空字符串表示不过滤
WindowSec int // 时间窗口(秒)
Threshold int // 触发阈值
}
// ruleCache 以 channelID 为 Key 的本地规则缓存
type ruleCache struct {
sync.RWMutex
data map[uint64][]*CachedRule
}
var cache = &ruleCache{
data: make(map[uint64][]*CachedRule),
}
为什么把 *regexp.Regexp 缓存起来?
regexp.Compile 是有开销的,如果每条日志都重新编译一遍正则,CPU 消耗非常可观。缓存预编译好的 *regexp.Regexp 句柄,匹配时直接调用 r.Regexp.MatchString(),性能提升显著——而且 *regexp.Regexp 是并发安全的,多个 goroutine 同时读没有问题。
查规则时,先走缓存,未命中再回源数据库:
go
func getRulesFromCache(ctx context.Context, channelID uint64) ([]*CachedRule, error) {
// 先加读锁查缓存
cache.RLock()
rules, exists := cache.data[channelID]
cache.RUnlock()
if exists {
return rules, nil // 命中缓存,直接返回
}
// 缓存未命中,查数据库(此处替换为你的查询逻辑)
dbRules, err := queryEnabledRules(ctx, channelID)
if err != nil {
return nil, err
}
// 预编译正则,组装 CachedRule
cached := make([]*CachedRule, 0, len(dbRules))
for _, r := range dbRules {
reg, err := regexp.Compile(r.Pattern)
if err != nil {
// 正则编译失败,记录日志跳过,不能因为一条脏数据崩掉整个引擎
log.Error("正则编译失败", "rule_id", r.ID, "pattern", r.Pattern, "err", err)
continue
}
cached = append(cached, &CachedRule{
ID: r.ID,
Name: r.Name,
Pattern: r.Pattern,
Regexp: reg,
Source: r.Source,
WindowSec: r.WindowSec,
Threshold: r.Threshold,
})
}
// 写入缓存(加写锁)
cache.Lock()
cache.data[channelID] = cached
cache.Unlock()
return cached, nil
}
规则变更时的缓存失效
当用户新增、修改、删除了一条告警规则,对应频道的缓存就必须失效,否则新规则不会生效、老规则不会退出。在规则的 CRUD Service 里调用这个函数即可:
go
// InvalidateCache 主动失效某个频道的规则缓存(在规则 CRUD 后调用)
func InvalidateCache(channelID uint64) {
cache.Lock()
defer cache.Unlock()
delete(cache.data, channelID)
}
五、把它们串在一起
整个告警引擎的判定流程:
go
// checkAlert 实时告警核心逻辑(在日志处理 Worker 中异步拉起,不阻塞主流程)
func checkAlert(ctx context.Context, channelID uint64, source, message string) {
// 1. 从缓存获取该频道的所有启用规则
rules, err := getRulesFromCache(ctx, channelID)
if err != nil {
log.Error("获取告警规则失败", "err", err)
return
}
for _, r := range rules {
// 2. 来源过滤(规则可以配置只匹配特定来源的日志)
if r.Source != "" && r.Source != source {
continue
}
// 3. 正则匹配(直接使用预编译好的句柄)
if r.Regexp == nil || !r.Regexp.MatchString(message) {
continue
}
// 4. 滑动窗口计数:在时间窗口内这条规则被命中了多少次?
count, err := SlideWindowCount(ctx, r.ID, r.WindowSec)
if err != nil {
log.Error("滑动窗口计数失败", "rule_id", r.ID, "err", err)
continue
}
// 5. 判断是否达到阈值
if count < int64(r.Threshold) {
continue // 还没到,继续等待
}
// 6. 检查冷却锁,防止告警风暴
if IsAlertCooldown(ctx, r.ID) {
continue // 还在冷却期,跳过
}
// 7. 设置冷却锁(3 分钟内不再重复告警)
_ = SetAlertCooldown(ctx, r.ID, 3*time.Minute)
// 8. 写入告警事件记录,并异步拉起 AI 根因分析
event := fireAlert(ctx, r, channelID, message, count)
go asyncAnalyzeAndBackfill(event.ID, channelID, r) // AI 分析,下一章详述
}
}
整个流程的数据流:
日志进来
│
▼
来源过滤(Source 字段)
│
▼
正则匹配(预编译 Regexp)── 不匹配 → 跳过
│
▼
SlideWindowCount (Redis ZSET)
│
▼
命中数 >= 阈值?── 否 → 跳过
│
▼
IsAlertCooldown (Redis Key)── 是 → 跳过(冷却中)
│
▼
SetAlertCooldown(锁定 3 分钟)
│
▼
写入 AlertEvent 记录
│
├── (同步)通知运维
│
└── (异步 goroutine)Gemini AI 根因分析 → 回填诊断报告
六、AI 根因分析:告警触发后让大模型帮你定位问题
告警发出去了,运维开始排查。下一步通常是:翻日志、找时间段、猜是哪里出了问题。
这个过程完全可以让大模型来做。告警触发的那一刻,我们已经知道:哪条规则被触发了、时间窗口是多少、用什么正则匹配到了什么——把这些信息连同触发窗口内的原始日志一起喂给模型,它能直接给出一份专业的 SRE 根因分析报告。
6.1 整体设计
告警触发 → 写入 AlertEvent 记录(此时 ai_analysis 字段为空)→ 异步拉起 AI 分析协程 → 等待日志 flush → 查询触发窗口内的日志 → 调用 Gemini → 回填 ai_analysis 字段。
整个分析过程完全异步,不阻塞告警主流程,用户打开告警详情时,报告可能已经生成好了。
6.2 一个必须有的等待:time.Sleep(1 * time.Second)
go
func asyncAnalyzeAndBackfill(eventID int, channelID uint64, rule *CachedRule) {
// 必须等待 1 秒!
// 日志处理采用批量写入策略(攒够 N 条或每 500ms flush 一次)
// 如果 AI 分析协程立即查数据库,触发告警的那批日志可能还在内存 Buffer 里
// 等 1 秒确保这批日志已经落盘,查到的上下文才是完整的
time.Sleep(1 * time.Second)
ctx := context.Background() // 见下方解释,不能用请求的 ctx
// 查询时间窗口内最近 20 条日志作为 AI 分析的上下文
startTime := time.Now().Add(-time.Duration(rule.WindowSec) * time.Second)
triggerLogs, err := queryRecentLogs(ctx, channelID, startTime, 20)
if err != nil {
log.Error("AI 分析:查询上下文日志失败", "err", err)
return
}
// 调用 Gemini 生成根因分析报告
report, err := analyzeWithGemini(ctx, rule, triggerLogs)
if err != nil {
log.Error("AI 分析:生成报告失败", "err", err)
return
}
// 回填到 AlertEvent 记录
if err := backfillReport(ctx, eventID, report); err != nil {
log.Error("AI 分析:回填失败", "err", err)
}
}
为什么用 context.Background() 而不是传入的 ctx?
告警引擎是从日志处理 Worker 里异步拉起的。如果这个 Worker 的 ctx 和某个 HTTP 请求或上报通道绑定,请求结束时 ctx 就会被取消,AI 分析协程的数据库查询和 Gemini 调用会被立即中断。
AI 分析调用 Gemini 可能需要几秒,必须用独立的 context.Background() 保证协程的生命周期不受外部影响。
6.3 为什么只取 20 条日志?
两个原因:
精准度:触发告警的关键日志通常就密集出现在时间窗口末尾,取太多反而把模型的注意力分散到无关日志上,分析结论会变模糊。20 条对根因定位来说已经足够。
Token 成本:每次 AI 调用都是真金白银。如果不限制日志数量,一次高频告警可能把几百条日志全塞进 prompt,Token 费用会非常可观。限制 20 条是在分析质量和成本之间取的平衡点。
6.4 Prompt 工程:让模型以 SRE 视角分析
模型输出的质量,90% 取决于 prompt 写得好不好。对于日志根因分析,推荐这样构建 prompt:
go
func buildPrompt(rule *CachedRule, logs []LogEntry) string {
var logText strings.Builder
for i, l := range logs {
logText.WriteString(fmt.Sprintf(
"[%d] [%s] [%s] %s | source: %s\n",
i+1, l.CreatedAt.Format("15:04:05"), l.Level, l.Message, l.Source,
))
}
return fmt.Sprintf(`你是一名经验丰富的 SRE 工程师,请分析以下告警事件的根本原因并给出处置建议。
## 告警规则
- 规则名称:%s
- 匹配模式(正则):%s
- 触发条件:%d 秒内匹配 %d 次
## 触发窗口内的原始日志(最近 %d 条)
%s
## 请输出以下内容(Markdown 格式):
1. **根本原因**:用一句话概括
2. **详细分析**:结合日志内容展开,指出关键错误模式
3. **影响范围**:推断受影响的服务或组件
4. **处置建议**:给出 2-3 条具体可操作的排查步骤`,
rule.Name, rule.Pattern, rule.WindowSec, rule.Threshold, len(logs), logText.String(),
)
}
6.5 Temperature = 0.2:调低随机性
Gemini 的 Temperature 参数默认值是 1.0。对于创意写作,高随机性是好事;但对于日志根因分析,我们要的是稳定、准确的技术输出,不需要模型发挥想象力。
设置 Temperature = 0.2 让模型更专注于日志内容,减少"幻觉(Hallucination)",输出更贴近实际问题:
go
modelInstance := client.GenerativeModel("gemini-2.0-flash")
var temperature float32 = 0.2
modelInstance.Temperature = &temperature
resp, err := modelInstance.GenerateContent(ctx, genai.Text(prompt))
6.6 实际效果
告警触发后大约 3-5 秒,AlertEvent 记录里的 ai_analysis 字段就会填入一份 Markdown 报告:
markdown
## 根本原因
数据库连接池耗尽,导致所有需要数据库操作的请求超时失败。
## 详细分析
从日志时间线来看,错误首次出现在 01:03:21,初始错误为 `connection refused`,
随后在 12 秒内迅速蔓延至全部 Worker...
## 影响范围
受影响服务:user-service、order-service
错误高峰期:01:03:21 - 01:05:47
## 处置建议
1. 立即检查数据库连接池配置,适当增大 max_open_conns
2. 查看数据库服务器负载,确认是否存在慢查询堆积
3. 检查最近是否有流量突增或批量任务并发执行
运维打开告警详情的时候,报告已经准备好了,可以直接按建议排查,不需要再手动翻日志。
七、几个容易被忽略的细节
1. 为什么 checkAlert 要异步调用?
告警引擎里有数据库查询和多次 Redis 操作,是有 I/O 开销的。如果同步执行,会阻塞日志处理的主流程,影响日志写入的吞吐量。用 go checkAlert(...) 异步拉起,日志处理主流程立即返回,告警判定在后台并发执行。
代价是:告警判定的时序和日志实际到达的时序可能有轻微偏差,但对于告警场景来说,几十毫秒的偏差完全可以接受。
2. 正则编译失败时,为什么要跳过而不是报错退出?
如果数据库里有一条历史脏数据(正则写错了),编译失败时直接 panic 或者返回 error 中断整个流程,会导致这个频道的所有规则都失效。
正确的做法是跳过这条坏规则、记录日志告警、继续处理剩余规则。让一条烂规则影响全局是不可接受的。
3. ZSET Key 为什么不需要单独设过期时间?
每次 SlideWindowCount 都会用 ZRemRangeByScore 清理过期的旧数据。在没有命中时,Key 里没有任何元素(或者元素被清空),这个 Key 对内存的占用就是零(Redis 的空 ZSET 会自动被 GC)。所以不需要额外的过期时间管理。
总结
| 组件 | 实现方式 | 解决的问题 |
|---|---|---|
| 滑动窗口计数 | Redis ZSET(Score = 时间戳,UUID Member) | 固定窗口的边界效应,精准统计任意时间段内的命中次数 |
| 内存规则缓存 | sync.RWMutex + map + 预编译正则 |
日志热路径不查 DB,规则变更时主动失效 |
| 告警冷却锁 | Redis Key + TTL | 防止同一规则在故障持续期间无限重复告警 |
| Redis 降级策略 | 故障时放行 | Redis 挂了不能把真实告警也屏蔽掉 |
| AI 根因分析 | Gemini + 异步协程 + 回填 | 告警触发后自动生成诊断报告,减少人工排查时间 |
整套系统的核心理念是分层处理:正则做第一层过滤(极低开销)→ 滑动窗口精准计数(防误报)→ 冷却锁降噪(防风暴)→ AI 智能诊断(提效)。每一层各司其职,可以根据需要单独替换或升级。
希望对你有帮助,有问题欢迎在评论区讨论!