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

7.4 KiB
Raw Blame History

tags, create time
tags create time
MQ
消费积压
性能优化
Kafka
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 逻辑,跳过非关键消息,优先处理核心业务。

// 降级消费:跳过非关键消息
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

关联笔记