消息堆积与依赖超时
消息堆积本身不是故障,堆积背后的“消费能力小于生产速率”才是。真正危险的是堆积之后的两件事:消费者被下游拖住导致线程池占满,以及重试把一次超时放大成一次雪崩。本章按成因、判读、处置三段展开,RabbitMQ 与 Kafka 的判读点分开讲。
三种成因,先分清
| 成因 | 特征 | 快速验证 |
|---|---|---|
| 生产速率突增 | 生产曲线是正常业务高峰形态,消费速率平稳 | 与昨日同时段生产量对比 |
| 消费能力下降 | 生产量不变,消费速率明显掉下来 | 消费者实例数、消费者错误日志 |
| 消费者掉线或重平衡 | 消费速率归零或剧烈波动,队列长度阶梯式上涨 | 消费者进程状态、rebalance 记录 |
分不清成因就会开错药:生产突增要限流或扩容,消费能力下降要先解决依赖问题,而重平衡频繁时扩容只会让情况更糟。
RabbitMQ 侧
rabbitmqctl list_queues name messages messages_ready messages_unacknowledged consumers
rabbitmqctl list_queues name messages consumers -p /prod
rabbitmqctl list_consumers
rabbitmqctl list_connections name state channels
rabbitmqctl status | grep -A5 -E 'memory|disk_free|file_descriptors'
| 指标 | 正常 | 异常含义 |
|---|---|---|
messages_ready | 稳定或小幅波动 | 持续增长:消费跟不上生产,或消费者被流控 |
messages_unacknowledged | 少量,随处理完成回落 | 一直高且不降:消费端卡在依赖上,或 prefetch 设置过大 |
consumers | 等于部署的实例数乘线程数 | 为 0 说明消费者全部掉线,堆积必然增长 |
| 连接 state | running / flowing | 出现 blocked:触发内存或磁盘流控,生产者发消息会被阻塞,应用侧表现为发送超时 |
| 队列 memory | 与消息量匹配 | 持续增长:消息体过大或大量消息长期未确认 |
| disk_free | 高于 disk_free_limit | 低于限制会同时阻塞生产者,先清磁盘再谈堆积 |
流控的两个开关:
rabbitmqctl set_vm_memory_high_watermark 0.6
rabbitmqctl set_disk_free_limit 10GB
内存或磁盘水位触发时,RabbitMQ 会阻塞生产者连接(而不是丢消息),所以上游看到的是“发消息超时”,不要误判成网络问题。
Kafka 侧
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 --describe --group order-group
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 --list
kafka-topics.sh --bootstrap-server kafka1:9092 --describe --topic order-event
| 字段 | 正常 | 异常含义 |
|---|---|---|
LAG | 稳定在低位 | 持续增长:消费能力不足 |
CURRENT-OFFSET | 持续推进 | 有流量但长时间不动:消费者卡住或 poll 间隔超时被踢出组 |
| 各分区 LAG | 各分区接近 | 单分区明显偏高:key 设计导致数据倾斜,或分区数不足 |
CONSUMER-ID 为 - | 有活跃消费者 | 消费者已离线,只剩位点信息 |
| rebalance 频繁发生 | 长时间稳定 | max.poll.interval.ms 太小或单批处理太慢,消费者被反复踢出组 |
关键取舍:
max.poll.interval.ms默认 300000(5 分钟),表示两次 poll 之间允许的最大间隔。单批处理超过它,消费者会被踢出组并触发 rebalance,这批消息会被重复消费。要么调大它,要么调小max.poll.records让单批更快处理完。- 重复消费无法彻底避免,消费端必须幂等:业务唯一键、状态机前置判断、去重表三者至少用一种。
- 分区数是消费并行度的上限:消费者实例数超过分区数时,多出来的实例只能空闲,此时扩容无效,要先扩分区(注意扩分区会改变 key 的分区映射)。
依赖超时如何级联
下游依赖变慢 → 消费线程被占用 → 消费速率下降 → 队列堆积
↓ ↓
请求等待与重试 → 线程池/连接池耗尽 → 接口全线超时 → 用户可感知故障
| 放大环节 | 后果 | 应对 |
|---|---|---|
| 调用无超时 | 线程被永久占用,池子很快耗尽 | 每个依赖必须设超时,且各级超时之和小于接口超时预算 |
| 无退避重试 | 重试流量是正常流量的 2 ~ 3 倍 | 有限重试加指数退避与随机抖动 |
| 同步调用非核心依赖 | 非核心故障拖垮核心链路 | 隔离舱(独立线程池与连接池)加熔断降级 |
| 消费者与在线服务共用资源 | 消费把在线接口拖慢 | 消费者独立部署、独立线程池与连接池 |
| 没有死信兜底 | 坏消息反复重投,永远消费不掉 | 重试上限 + 死信队列 + 定时补偿任务 |
消费能力估算与扩容决策
预计消化时间 = 当前积压量 /(消费速率 - 生产速率)
消费速率 ≈ 单消费者吞吐 × 消费者数 × 批量修正系数
| 场景 | 扩容是否有效 | 说明 |
|---|---|---|
| Kafka 消费者数少于分区数 | 有效 | 一直加到等于分区数为止 |
| Kafka 消费者数已等于分区数 | 无效 | 先扩分区,或提升单条消息的处理速度 |
| RabbitMQ 竞争消费 | 有效 | 但要看 broker 的 CPU 与网络是否成为新瓶颈 |
| 消费逻辑被外部依赖拖住 | 有限 | 先做超时与隔离,否则加实例只会把依赖压死 |
| 消息体过大导致网络瓶颈 | 无效 | 先压缩或拆小消息体 |
判断扩容有没有用,只用一个标准:瓶颈是否在消费者自身。瓶颈在下游依赖或 broker 上时,扩容只会让更多人排队。
常见误区
| 误区 | 后果 | 正确做法 |
|---|---|---|
| 只按队列长度绝对值告警 | 短时高峰被误报,真故障反而被淹没 | 用预计消化时间与速率曲线告警 |
| 重试没有上限 | 坏消息反复占用消费能力 | 重试次数上限加死信队列 |
| 消费端不做幂等 | 重试与重平衡造成重复业务数据 | 唯一键、去重表或状态机 |
| 盲目增大 prefetch 提升吞吐 | 未确认消息堆积,内存上涨且故障时大量重投 | 按单条处理耗时设置合理的 prefetch |
处置动作
- 先止损再定位:临时扩容消费者;受分区数限制的场景扩容无效,改为增加消费线程或批量大小。
- 生产端限流:对非核心消息直接限速或暂停生产,把消费能力让给订单、支付类消息。
- 非核心降级:营销、统计、积分类消息暂停消费,等水位降下来再补。
- 卡住的消费者:确认是依赖超时还是死循环,先重启恢复消费能力,同时保留日志与线程栈。
- 坏消息直投死信:格式错误或必然失败的消息不要留在主队列反复阻塞分区。
- 谨慎操作:
rabbitmqctl purge_queue清空队列不可逆,除非业务确认消息可丢,否则一律走死信与补偿。
预防与巡检
- 同时监控三条曲线:生产速率、消费速率、队列长度(Lag)。Lag 告警按“预计消化时间”而不是绝对值——Lag 10 万但一分钟能消化完,不等于故障。
- 消费端把幂等、超时预算、重试次数、死信策略写进脚手架,靠规范约束而不是靠人记得。
- 大促前做消费能力压测,明确单消费者 TPS 与水平扩展上限,提前扩分区。
- 健康检查要看消费速率而不是进程存活:进程活着但完全不消费,是最隐蔽的堆积原因。
小结:堆积先分清是生产突增、消费变慢还是消费者掉线,RabbitMQ 盯 messages_ready、未确认数与连接 blocked 状态,Kafka 盯 LAG、分区倾斜与 rebalance;真正要防的是依赖超时通过线程池与重试放大成雪崩,因此超时预算、熔断隔离、幂等消费与死信兜底必须提前建设。