7.4 KiB
tags, create time
| tags | create time | ||||
|---|---|---|---|---|---|
|
2026-05-24 19:52 |
MQ 消费积压治理
概述
消息积压是 MQ 使用中最常见的"事故"之一:Producer 疯狂写入,Consumer 跟不上节奏,Lag 越堆越高,下游业务开始告警。本文从积压的成因、紧急处理到根本解决方案,系统性地梳理治理思路。
正文
消息积压是怎么产生的
消息积压的本质很简单:生产速度 > 消费速度。但具体原因千差万别:
- Consumer 消费速度跟不上:业务逻辑变重了(比如新增了一次远程调用),单条消息处理耗时从 10ms 涨到 500ms,消费速率直接腰斩。
- Consumer 故障:Consumer 实例挂了、OOM 了、被 K8s 驱逐了,但你不知道——因为消息队列不会主动告诉你"你的消费者已经死了"。
- Rebalance 风暴:Consumer Group 频繁 Rebalance,每次 Rebalance 期间所有 Consumer 都暂停消费。如果 Consumer 处理超时触发 Rebalance,然后 Rebalance 又导致更多超时,就会进入恶性循环。
- 流量突增:大促、秒杀、爬虫攻击,Producer 端流量瞬间翻几倍,但 Consumer 还是原来那点资源。
积压的影响
积压不是"慢一点"那么简单,它会引发连锁反应:
- 消息延迟增大:用户下单后要等几分钟才能收到确认通知,体验极差。
- 触发过期删除:Kafka 的日志保留策略(比如 7 天),如果积压的消息超过了保留时间,还没被消费就被删了,数据丢失。
- 下游业务受影响:如果下游依赖 CDC 事件更新缓存或索引,积压会导致缓存长时间不更新,数据不一致。
紧急处理方案
发现积压后,第一步不是优化代码,而是先止血。
方案一:快速扩容 Consumer 实例
这是最直接的办法。Consumer Lag 了 10 万条?加 Consumer 实例就行。但有个前提:Partition 数量要够。
Kafka 中,一个 Partition 最多只能被同一个 Consumer Group 中的一个 Consumer 消费。如果你的 Topic 只有 4 个 Partition,那最多只能有 4 个 Consumer 同时消费,加再多实例也没用。
[!question] 你的 Topic 有 8 个 Partition,当前跑了 3 个 Consumer 实例,Lag 在持续增长。这时候你决定加到 10 个 Consumer,会发生什么?多出来的 2 个在干嘛?
方案二:临时 Topic 转发
如果 Consumer 的消费逻辑很重(比如涉及数据库写入、远程调用),而且短时间内改不了,可以这样做:
- 写一个"轻量消费者",只做一件事:从积压的 Topic 读消息,批量转发到一个临时 Topic。
- 启动一批新的 Consumer 消费临时 Topic,这些 Consumer 的消费逻辑可以是更轻量的版本(比如只写入一个临时表,不做完整业务处理)。
- 积压消化完后,再慢慢回补业务逻辑。
这个方案的核心思想是:先快速消费完,再慢慢处理。
方案三:降级消费
如果积压的消息中有关键消息和非关键消息(比如订单消息是关键的,日志消息是非关键的),可以临时修改 Consumer 逻辑,跳过非关键消息,优先处理核心业务。
// 降级消费:跳过非关键消息
func handleMsg(msg *kafka.Message) {
priority := msg.Headers.Get("priority")
if priority == "low" && isBacklogHigh() {
// 积压严重时,低优先级消息直接丢弃
log.Printf("skip low priority msg: offset=%d", msg.Offset)
return
}
// 正常处理关键消息
processBusinessLogic(msg)
}
[!question] 降级消费意味着丢弃部分消息。在什么业务场景下这是可以接受的?你能想到哪些场景绝对不能降级?
根本解决:优化消费逻辑
紧急处理只是止血,要根治得从消费逻辑入手。
减少 IO 操作:消费逻辑中的数据库写入、远程调用是最常见的瓶颈。检查一下是不是每条消息都触发一次 DB 写入?能不能改成批量写入?
异步处理:消费逻辑中如果有非关键步骤(比如发通知、写日志),可以异步化。消费者只做最关键的操作(比如写主库),其余的丢到 goroutine 或者另一个队列。
批量消费:攒一批消息一起处理,减少网络往返和事务开销。
package main
import (
"context"
"log"
"time"
"github.com/segmentio/kafka-go"
)
func main() {
reader := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"localhost:9092"},
Topic: "orders",
GroupID: "order-processor",
MinBytes: 1, // 最少 1 字节就返回
MaxBytes: 10e6, // 最多 10MB
CommitInterval: 0, // 手动提交
})
defer reader.Close()
batch := make([]kafka.Message, 0, 100)
ticker := time.NewTicker(500 * time.Millisecond) // 每 500ms 刷一次
defer ticker.Stop()
for {
select {
case <-ticker.C:
if len(batch) == 0 {
continue
}
// 批量写入数据库(一个事务搞定)
if err := batchInsert(batch); err != nil {
log.Printf("batch insert error: %v", err)
continue
}
// 手动提交 offset
if err := reader.CommitMessages(context.Background(), batch[len(batch)-1]); err != nil {
log.Printf("commit error: %v", err)
}
log.Printf("processed batch of %d messages", len(batch))
batch = batch[:0] // 清空批次
default:
ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond)
msg, err := reader.ReadMessage(ctx)
cancel()
if err != nil {
continue // 超时,继续攒消息
}
batch = append(batch, msg)
if len(batch) >= 100 {
// 批次满了,立即处理
ticker.Reset(0) // 触发立即刷入
}
}
}
}
func batchInsert(msgs []kafka.Message) error {
// 批量插入数据库的逻辑
return nil
}
这段代码展示了批量消费的核心思路:攒够 100 条或者等 500ms,然后一次性写入数据库。相比逐条消费逐条写入,批量处理可以将数据库写入次数降低一到两个数量级。
增加 Partition:如果优化消费逻辑后还是不够,那就增加 Partition 数量,从根本上提升消费并行度。但这需要新建 Topic、迁移数据、修改 Producer 配置,成本不低,应该提前做好容量规划。
预防措施
与其事后救火,不如事前防火:
- 容量规划:根据历史流量峰值,预留 2-3 倍的消费能力。
- 压力测试:上线前用压测工具(如 Kafka 的
kafka-producer-perf-test)验证消费能力。 - 告警阈值:设置 Lag 告警,比如 Lag > 5000 触发 P2 告警,Lag > 50000 触发 P1 告警,Lag 连续 5 分钟单调递增触发紧急告警。
积压处理决策流程
flowchart TD
A["发现消费积压"] --> B{"积压原因?"}
B -->|"Consumer 实例挂了"| C["重启/扩容 Consumer"]
B -->|"消费逻辑太慢"| D{"能否快速优化?"}
B -->|"流量突增"| E{"Partition 数量够?"}
D -->|"能"| F["异步化/批量消费"]
D -->|"不能"| G["临时 Topic 转发"]
E -->|"够"| C
E -->|"不够"| H["降级消费 + 后续扩容 Partition"]
C --> I["观察 Lag 趋势"]
F --> I
G --> I
H --> I
I -->|"Lag 恢复"| J["复盘 + 优化预案"]
I -->|"Lag 未恢复"| B