做实时日志系统的时候,告警功能几乎是标配。 最开始我的实现很简单粗暴:日志进来,逐条匹配规则,满足
用 Go + Redis ZSET 实现滑动窗口告警引擎与 Gemini AI 根因分析
发布时间: 2026-07-17 (5 days ago)
GOGemini

做实时日志系统的时候,告警功能几乎是标配。

最开始我的实现很简单粗暴:日志进来,逐条匹配规则,满足条件就发告警。

结果线上跑了一段时间,有一次某个服务疯狂报错,告警系统在两分钟内发出了 200 多条通知。运维的手机被叮个不停,最后直接把通知关掉了——反而错过了真正需要关注的窗口期。

这就是告警风暴(Alert Storm),也是大部分自研告警引擎最容易踩的坑。

这篇文章讲三件事:

  1. 用 Redis ZSET 实现滑动窗口计数,取代"每条日志都触发告警"的简单计数
  2. 用 Redis 冷却锁防止同一规则在短时间内重复告警
  3. 告警触发后,异步调用 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 智能诊断(提效)。每一层各司其职,可以根据需要单独替换或升级。

希望对你有帮助,有问题欢迎在评论区讨论!