面试知识库
高 进阶

消息积压处理方案#

一句话答案#

紧急处理:扩消费者实例 + 临时调大消费线程 + 跳过非核心消息;根治:优化消费逻辑/合理设置 Partition 数。

核心要点

紧急处理:

  1. 扩容消费者实例(Kafka 经典消费者组需先扩Partition;Kafka 4.2 起生产可用的 Share Group(KIP-932)和 RocketMQ 5.x POP 消费允许消费者数超过分区/队列数,但不保证顺序)
  2. 临时调大消费线程数
  3. 跳过/转存非核心消息
  4. 新建临时Topic分流

根治: 优化消费逻辑(批量处理/异步化) / 合理设置Partition数

面试回答(2分钟版)

消息积压本质是消费速度跟不上生产速度,处理分紧急和根治两个层面。紧急处理首先是扩容消费者实例来提高消费并行度,但 Kafka 要注意消费者数不能超过 Partition 数,所以先把消费者补到和 Partition 数一样多;注意临时扩 Partition 救不了存量积压——新分区只接新消息,老消息还在原分区;其次可以临时调大单个消费者的消费线程数;对于非核心消息可以先跳过或转存到其他 Topic 后续慢慢处理,优先保证核心业务消息的消费;还可以新建临时 Topic 把积压消息分流到更多消费者处理。根治层面要从消费逻辑本身入手:优化消费处理逻辑减少单条消息的处理时间,比如批量处理替代逐条处理、异步化 IO 操作、减少不必要的远程调用;合理设置 Partition 数量让消费者能充分并行。另外还要建立监控告警机制,对消费延迟设置阈值,积压达到一定量就自动告警,避免问题积累到不可控才发现。

追问与易错

追问方向:

  • “消息积压的常见原因有哪些?”→ 消费者处理慢(下游服务超时/GC/慢 SQL)、消费者实例数不足或挂了、生产端突发流量、Rebalance 频繁导致消费停顿
  • “积压了几百万消息怎么紧急处理?”→ 临时扩容:新建一个多分区的临时 Topic,写一个只做转发的消费者把积压消息分发到新 Topic,然后用大量消费者并行消费临时 Topic 快速消化
  • “怎么预防消息积压?”→ 提前做容量规划(分区数≥消费者数)、消费端异步化批量处理、设置消费延迟监控告警(如 Kafka lag 监控)、消费者做好限流降级避免雪崩

易错点:

  • ❌ 积压了就加消费者实例——Kafka 同一消费组里超过分区数的消费者会空闲(4.2 起的共享组除外);先看分区数,不够再扩分区或用临时 Topic 转发扩容
  • ❌ 积压时调大 max.poll.records 多拉一些——单批处理时间一旦超过 max.poll.interval.ms(默认 5 分钟),消费者会被踢出组、触发 Rebalance,反而更慢
  • ❌ 只清积压不查根因——先分清是消费端变慢(下游慢 SQL、外部接口超时、异常重试)还是生产端流量突增,否则扩容后很快又会积压