Files
cs-note/hhs/MQ/10-监控与运维/37-MQ-消费积压治理.md
T
2026-05-24 20:51:06 +08:00

193 lines
7.4 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
tags:
- MQ
- 消费积压
- 性能优化
- Kafka
create time: 2026-05-24 19:52
---
# MQ 消费积压治理
## 概述
消息积压是 MQ 使用中最常见的"事故"之一:Producer 疯狂写入,Consumer 跟不上节奏,Lag 越堆越高,下游业务开始告警。本文从积压的成因、紧急处理到根本解决方案,系统性地梳理治理思路。
## 正文
### 消息积压是怎么产生的
消息积压的本质很简单:**生产速度 > 消费速度**。但具体原因千差万别:
1. **Consumer 消费速度跟不上**:业务逻辑变重了(比如新增了一次远程调用),单条消息处理耗时从 10ms 涨到 500ms,消费速率直接腰斩。
2. **Consumer 故障**:Consumer 实例挂了、OOM 了、被 K8s 驱逐了,但你不知道——因为消息队列不会主动告诉你"你的消费者已经死了"。
3. **Rebalance 风暴**:Consumer Group 频繁 Rebalance,每次 Rebalance 期间所有 Consumer 都暂停消费。如果 Consumer 处理超时触发 Rebalance,然后 Rebalance 又导致更多超时,就会进入恶性循环。
4. **流量突增**:大促、秒杀、爬虫攻击,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 的消费逻辑很重(比如涉及数据库写入、远程调用),而且短时间内改不了,可以这样做:
1. 写一个"轻量消费者",只做一件事:从积压的 Topic 读消息,批量转发到一个临时 Topic。
2. 启动一批新的 Consumer 消费临时 Topic,这些 Consumer 的消费逻辑可以是更轻量的版本(比如只写入一个临时表,不做完整业务处理)。
3. 积压消化完后,再慢慢回补业务逻辑。
这个方案的核心思想是:**先快速消费完,再慢慢处理**。
**方案三:降级消费**
如果积压的消息中有关键消息和非关键消息(比如订单消息是关键的,日志消息是非关键的),可以临时修改 Consumer 逻辑,跳过非关键消息,优先处理核心业务。
```go
// 降级消费:跳过非关键消息
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 或者另一个队列。
**批量消费**:攒一批消息一起处理,减少网络往返和事务开销。
```go
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 分钟单调递增触发紧急告警。
### 积压处理决策流程
```mermaid
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
```
## 关联笔记
- [[30-MQ-核心概念与选型]]
- [[33-MQ-Kafka-架构与核心机制]]
- [[36-MQ-监控指标与告警]]
- [[35-MQ-与-CDC]]