← 返回题目列表

消息队列出现消息积压怎么办?

高频 中等 第 3 / 25 题 更新于 2026/07/28
消息队列消息积压消费能力扩容

简化版

消息积压指生产速度远大于消费速度,导致大量消息堆积在队列里、消费延迟越来越大。处理分应急根治两方面:应急——快速提升消费能力:给消费者组扩容(增加消费者实例,但受分区数限制)、或先用一个消费者把积压消息快速转发到一个临时扩容的新 Topic(更多分区)再多消费者并行处理。根治——排查积压根因:是消费者代码慢(优化逻辑、批量处理、异步化)、还是下游依赖慢(DB/接口瓶颈)、还是分区数不够(增加分区提升并行度)、或消费者故障。

详细版

积压的常见原因:

  • 消费者处理逻辑太慢(复杂计算、慢查询、调用慢的下游)。
  • 消费者数量不够 / 分区数太少,并行度上不去。
  • 消费者出 bug 卡住、宕机、频繁 Rebalance。
  • 突发流量导致生产远超消费。

应急处理(快速消化积压):

方案一:直接扩容消费者
  - 增加消费者实例(Kafka 中消费者数 ≤ 分区数才有效)

方案二:临时扩容(分区不够时)
  1. 加一个消费者,只负责把积压消息快速搬运到新建的、分区更多的临时 Topic
  2. 对新 Topic 启动大量消费者并行处理
  3. 消化完后恢复原架构

根治手段:

  • 优化消费逻辑:批量拉取、批量处理、批量写库;耗时操作异步化。
  • 增加分区数 + 相应增加消费者,提升并行度。
  • 提升下游能力:DB 加索引、加缓存、扩容。
  • 加监控告警:对积压量设阈值,及时发现。

完整版教学

一、什么是积压:生产 > 消费的必然结果

MQ 里,生产者往队列写、消费者从队列读。如果生产速度持续大于消费速度,队列里的消息就会越堆越多,这就是消息积压(堆积、Lag)。表现是消费延迟越来越大——一条消息可能要等几分钟、几小时才被处理,严重时磁盘被撑满、MQ 崩溃。

积压本身是「削峰」的正常现象(短时高峰后能慢慢消化就没事),但如果积压持续增长、消化不掉,就是问题,说明消费能力长期跟不上。

二、先定位原因,再对症下药

处理积压不能盲目扩容,要先搞清为什么积压

  • 消费者处理太慢:消费逻辑里有复杂计算、慢 SQL、调用了慢的外部接口。
  • 并行度不够:消费者实例太少,或分区数太少(Kafka 里消费者并行度上限 = 分区数,分区不够加消费者也没用)。
  • 消费者异常:消费者 bug 卡住、频繁抛异常重试、宕机、频繁 Rebalance 导致消费停滞。
  • 突发流量:生产端突然暴增(大促、故障重放),瞬间远超消费能力。

不同原因解法不同——是代码慢就优化代码,是分区不够就加分区,是流量暴增就临时扩容。

三、应急处理:快速消化积压

线上已经严重积压、要尽快消化时:

方案一:直接扩容消费者。 如果分区数还够(消费者数 < 分区数),直接增加消费者实例,让更多消费者并行消费不同分区。注意 Kafka 里一个分区同一时刻只能被消费者组里的一个消费者消费,所以消费者数超过分区数是无效的(多出的消费者空闲)。

方案二:临时扩容(分区不够时的经典手段)。 如果分区数本身就少(比如只有 4 个分区,加消费者也只能 4 个并行),无法快速加分区消化时:

  1. 新建一个分区数很多的临时 Topic;
  2. 写一个轻量的中转消费者,它不做业务处理,只负责把积压 Topic 的消息快速搬运到新的临时 Topic(搬运很快,因为不做重活);
  3. 对临时 Topic 启动大量消费者并行处理业务;
  4. 积压消化完后,恢复原来的架构。

这个「先搬运扩分区、再并行消费」是应对突发积压的经典套路。

四、根治:从根本上提升消费能力

应急之后要根治,让消费能力长期匹配生产:

① 优化消费逻辑(最常见):

  • 批量处理:一次拉取多条消息、批量写库/批量调用(减少 I/O 次数)。
  • 异步化:把消费逻辑里的耗时操作异步处理,缩短单条消息处理时间。
  • 去掉慢操作:优化慢 SQL(加索引)、加缓存、减少不必要的远程调用。

② 提升并行度:

  • 增加分区数 + 相应增加消费者。这是提升并行消费能力的根本(分区是并行的基本单位)。注意增加分区对已有的 key 顺序有影响,要评估。

③ 提升下游能力:

  • 消费慢往往是被下游拖累(DB 写不动、接口慢)。要给下游 DB 加索引、加缓存、扩容、做批量。

④ 加监控告警:

  • 对**消费 Lag(积压量)**设置监控告警,积压超过阈值就及时介入,别等到磁盘满了才发现。

五、预防:让积压不发生

比事后处理更重要的是预防:

  • 容量规划:根据峰值流量评估需要的分区数、消费者数,留足余量。
  • 压测:上线前压测消费能力,确保能扛住预期峰值。
  • 监控 Lag:常态化监控消费延迟,趋势异常提前处理。
  • 消费者健壮性:做好异常处理、限流、降级,避免消费者卡死。

六、常见误区与追问

考点正确口径
消息积压生产速度长期大于消费速度
排查维度生产突增、消费变慢、下游慢、分区不足、失败重试
处理扩消费者、扩分区、限流削峰、临时跳过或转储
lag = latest_offset - committed_offset
if lag grows continuously:
  check consumer error rate
  check downstream latency
  add consumers up to partition count

积压不是只看队列里有多少消息,而要看 lag 是否持续增长以及消费端为什么跟不上。

  • 误区:消息积压只要加消费者就能解决。 如果分区数不足、下游数据库慢或消费者一直报错,加消费者也没用。
  • 误区:积压时应该无限重试。 毒消息反复重试会卡住队列,应进入死信或隔离处理。
  • 误区:清空队列是首选方案。 清空会丢业务数据,只有确认可丢弃并做好审批审计才可做。
  • 追问:Kafka 消费者数量为什么受分区限制? 同一组内一个分区只能给一个消费者,消费者超过分区数会空闲。
  • 追问:如何快速止血? 限流生产、扩容消费者和下游、隔离慢逻辑、跳过毒消息或临时转储。
  • 追问:积压恢复后要复盘什么? 峰值生产速率、消费耗时、失败率、分区容量和告警阈值。

七、加强记忆

消息积压 = 生产速度持续 > 消费速度,消息堆积、消费延迟增大。处理先定位原因(消费逻辑慢 / 分区少并行不够 / 消费者异常 / 突发流量)。应急快速消化:分区够就直接加消费者(注意 Kafka 消费者数 ≤ 分区数才有效);分区不够用临时扩容套路——加中转消费者把积压消息快速搬到「多分区临时 Topic」,再启动大量消费者并行处理。根治:优化消费逻辑(批量、异步、去慢操作)、增加分区+消费者提并行、提升下游 DB/接口能力、加 Lag 监控告警。核心口诀:先看是消费慢还是并行少,应急扩容消化、根治提消费能力