先说结论

先确认堆积范围、增长速度和业务影响并停止流量放大,再区分生产突增、消费者异常、处理变慢、分区倾斜或下游瓶颈;短期通过恢复消费者、限流和增加有效并行度止损,长期通过容量规划、批量处理、背压和降级机制防止复发。

什么叫消息堆积

对某个消费者组,最新日志位点与已提交消费位点之间的差值就是 Lag。Lag 偶尔上升不一定异常,关键看它是否持续增长以及业务能接受的处理延迟。

积压增长速度 = 生产速率 - 消费速率
预计清空时间 = 当前积压量 / 恢复后的净消费速率

如果生产每秒 10 万条、消费每秒 8 万条,Lag 每秒增加 2 万。即使消费者仍在工作,系统也处于持续失血状态。

第一步:确认影响范围

收到告警后先回答:

  • 哪个集群、Topic、消费者组和分区堆积?
  • Lag 是均匀增长还是集中在少数分区?
  • 消费者实例是否存活,活跃成员数量是否变化?
  • 生产速率是否突增,消费速率是否下降?
  • 最老未消费消息已经延迟多久?
  • Topic 保留时间还剩多少安全窗口?
  • 哪些业务结果、通知或数据同步受到影响?

只看消费者组总 Lag 会掩盖分区倾斜。必须下钻到 Partition。

第二步:先止损

1. 停止重试风暴

下游超时后立即无限重试,会让同一批消息反复占用线程。应使用有限重试、指数退避,并把持续失败消息送入重试 Topic 或死信队列。

2. 限制上游非关键流量

若系统已经过载,可以暂停低优先级事件、降低采样率或让生产端缓冲。不能控制上游时,也要避免消费者因过载频繁崩溃和再均衡。

3. 隔离毒消息

某条格式错误或业务异常消息可能反复失败,阻塞整个分区。记录原消息、异常和重试次数,超过阈值后转入死信队列,让主消费链路继续推进。

资金等关键消息不能直接跳过,必须建立告警和人工修复流程。

常见根因一:消费者进程异常

检查消费者是否频繁重启、OOM、Full GC、线程池耗尽或发生未捕获异常。若实例不健康,单纯扩容只会增加更多失败实例。

同时检查再均衡频率。消费者处理超过 max.poll.interval.ms,会被移出组并反复再均衡,表现为实例都存活但实际吞吐很低。

常见根因二:下游依赖变慢

消费者通常要写 MySQL、Redis、ES 或调用接口。下游连接池耗尽、慢 SQL、锁等待和限流都会把消费线程阻塞。

应对方法:

  • 先恢复或降级下游服务。
  • 批量写数据库,减少网络往返。
  • 缩短事务并补充必要索引。
  • 设置合理超时,避免线程永久等待。
  • 对非关键下游采用异步或暂存。

如果下游每秒只能处理 1 万条,把 Kafka 消费者扩到每秒拉取 10 万条只会把压力转移并拖垮下游。

常见根因三:生产流量突增

活动、爬虫、批量回放或上游 Bug 可能让生产量突然翻倍。应确认是正常峰值还是异常流量,并判断容量是否足以在业务时限内追平。

正常峰值需要预先扩容和压测;异常生产需要限流、修复生产者并识别已经产生的无效消息。

常见根因四:分区倾斜

总消费能力充足,但某个热点 Key 把大量消息写入一个分区。该分区只能由组内一个消费者处理,其他消费者空闲。

常见原因:

  • 分区 Key 分布不均。
  • 大客户或超级租户占据单个 Key。
  • 空 Key 被统一路由。
  • 自定义分区器存在缺陷。

修改分区策略会影响同 Key 有序性。可以拆分热点 Key、为超级租户使用独立 Topic,或在业务允许时增加更细粒度的并行键。

常见根因五:分区数不足

消费者组内一个分区同时只分配给一个消费者。若有 12 个分区,增加到 30 个消费者,仍只有 12 个消费者工作。

增加分区可以提高并行上限,但要注意:

  • 旧 Key 到分区的映射可能变化。
  • 全局或 Key 有序语义需要重新评估。
  • Broker 元数据、文件句柄和复制开销增加。
  • 已经存在的积压不会自动均匀迁移到新分区。

因此临时增加分区并不一定能快速清掉旧积压。

常见根因六:单条消息处理太慢

消费者可能逐条查询、逐条写入,形成 N+1。可以在不破坏顺序和事务边界的前提下批量处理:

单条写数据库 1000 次
        ↓
按分区聚合后批量写 10 次

还可以缓存重复读取的维度数据、减少日志同步输出、预编译表达式和优化序列化。

如何临时扩容

横向增加消费者

当活跃消费者数少于分区数,且瓶颈在消费者 CPU 或本地处理时,增加实例通常有效。扩容后观察实际分区分配、消费速率和下游压力。

提高单实例并行度

poll 线程可把不同分区任务交给工作线程,但必须维护分区内顺序和位点水位。不能任务一提交线程池就立刻提交位点。

建立临时快速消费组

极端积压时,可以将消息转存到临时 Topic,使用更多分区和专用消费者快速处理。但这属于高风险方案,需要确保不丢失、不乱序、可对账,并限制临时链路重复处理。

消息快过期了怎么办

当最老消息接近 Topic 保留期限,优先扩大安全窗口:

  1. 评估并延长 Topic 保留时间。
  2. 确认 Broker 磁盘容量足够。
  3. 将未消费数据复制或导出到临时存储。
  4. 暂停可能加速磁盘耗尽的非关键流量。
  5. 建立消息 ID 清单和对账任务。

直接无限延长保留时间可能填满磁盘,引发更大故障。必须同时计算积压字节量、增长率和剩余空间。

恢复后如何安全追平

恢复消费能力后,不要立即解除所有限流。持续观察:

  • Lag 是否稳定下降。
  • 消费错误率和重试量。
  • MySQL、Redis、ES 等下游容量。
  • 消费者 GC、CPU、线程池和连接池。
  • Broker 网络、磁盘和请求延迟。

追平期间可能大量触发过期业务,例如延迟数小时的通知、优惠券或风控事件。消费者要判断消息时效,决定执行、丢弃还是补偿。

预计清空时间怎么计算

假设当前积压 3.6 亿条,生产速率 2 万条/秒,扩容后消费速率 5 万条/秒:

净清理速度 = 5 万 - 2 万 = 3 万条/秒
预计清空时间 = 3.6 亿 / 3 万 = 12000 秒 ≈ 3.3 小时

实际还要考虑流量波动、失败重试和下游限流。清空目标必须使用净消费速率,不能直接用总消费速率。

长期治理方案

  • 为每个消费者组设置 Lag、最老消息年龄和消费速率告警。
  • 压测峰值生产与消费能力,保留安全余量。
  • 为下游依赖设置超时、隔离、限流和降级。
  • 建立重试 Topic、死信队列和人工补偿平台。
  • 监控分区倾斜和消费者再均衡。
  • 规范批量回放和人工重置位点操作。
  • 为消息设置业务时效和过期处理策略。
  • 定期演练消费者全停、下游故障和大规模积压。

容易踩坑的地方

  • 看到 Lag 就无限增加消费者,忽略分区数和下游容量。
  • 为清积压直接提交最新位点,相当于主动丢弃所有历史消息。
  • 无限重试毒消息,导致一个分区永久阻塞。
  • 只看总 Lag,不检查单分区热点。
  • 追平后立即缩容,没有观察重试和业务补偿。
  • 延长保留时间却不计算磁盘剩余空间。

常见问题

追问 1:消息堆积后直接增加消费者有用吗?

只有消费者数少于分区数、瓶颈确实在消费者且下游还有容量时才有效。消费者超过分区数会空闲,下游已饱和时扩容反而加剧故障。

追问 2:某一个分区积压特别严重怎么办?

检查分区 Key 和热点业务。短期可为热点业务扩容单个处理链路或转移数据,长期要拆分热点 Key、调整分区策略或隔离超级租户,同时评估顺序语义。

追问 3:能不能重置位点到最新位置快速恢复?

这会跳过全部积压,相当于丢消息。只有业务明确允许丢弃、数据已备份且完成审批对账时才能使用,关键消息不能这样处理。

追问 4:如何判断多久可以清完?

使用当前积压量除以净清理速度,即扩容后消费速率减去同时发生的生产速率,并为重试和波动预留余量。

追问 5:消费者没有报错为什么还会积压?

可能是处理速度稳定但低于生产速度、某分区热点、下游响应变慢、批次过小或频繁再均衡。没有异常日志不代表吞吐满足需求。

追问 6:积压恢复后要做什么?

核对消息与业务结果、处理死信和过期事件,逐步解除限流,复盘容量和告警缺口,并通过压测与故障演练验证修复。