消息队列积压处理:先判断生产突增还是消费退化

通过入队速率、消费速率、消息年龄和失败分布诊断积压,并安全提升消费能力。

文章目录 · 4 节

队列长度上升只说明流入大于流出。应先判断是正常流量突增、消费者变慢、分区热点还是毒消息重试,否则盲目扩容可能放大下游故障。

观察四个核心信号

produce_rate
consume_rate
oldest_message_age
retry_or_dlq_rate

消息年龄比队列长度更接近用户影响。分区队列还要逐分区比较 lag,整体平均值可能掩盖单个 key 热点。

排查消费者

检查实例健康、处理耗时、错误率、rebalance、下游连接池和提交 offset 的方式。消费者日志要能关联 topic、partition、offset 与失败类别。

kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group order-worker

安全恢复

  • 下游健康且分区足够时再逐步增加消费者
  • 暂停无关生产者或降低非关键消息优先级
  • 对毒消息设置有限重试并进入死信队列
  • 不未经评估直接跳过 offset 或清空队列

收尾清单

积压清空后验证业务状态和重复消费影响。根据峰值输入、单条处理耗时和恢复目标计算所需余量,并对消息年龄设置 SLO 告警。