消费组重平衡,队列堆积了 30 万条
我们一个交易通知服务用 RocketMQ 4.9 做消费,某天上午扩容,新增 2 个消费者实例。按理说扩容应该更快,结果监控上"消费积压"从 0 飙到 30 万条,持续了 20 多分钟才消化完。组员在群里贴了张图:"加机器反而更慢了?"
根因在老架构上:RocketMQ 4.x 是队列级负载,一个 Topic 的 MessageQueue 在消费组内按平均分配。扩容触发重平衡,队列在实例间搬来搬去,搬运期间这些队列的消费是停摆的;再加上我们为了保顺序用了 ConsumeFromWhere 加锁,重平衡时锁竞争更重。
RocketMQ 5.0 的新架构:Proxy 与 gRPC
今年(2022)RocketMQ 5.0 把"计算"和"存储"拆开了,引入了一层 Proxy。Producer/Consumer 不再直连 Broker,而是先连 Proxy,由 Proxy 统一和 Broker 的存储层打交道。这层 Proxy 带来了两个直接好处:
- 多语言客户端有了统一入口。5.0 提供了标准 gRPC 客户端,Go、C++、Python 都能用同一套协议,以前各语言自己实现 remoting 协议是很痛的事。
- 客户端的负载均衡逻辑下沉到 Proxy/服务端,客户端变成"无状态"的,重平衡的代价大幅下降。
# 5.0 的 gRPC 客户端,Endpoint 指向 Proxy
Producer producer = Provider.getProducerBuilder()
.setClientConfiguration(
ClientConfiguration.newBuilder()
.setEndpoints("rmq-proxy-1:8081")
.build())
.setTopics("trade-notify")
.build();
Pop 消费:专为"队列堆积"设计
5.0 的 PopConsumer 是最让我感兴趣的改动。传统 Pull/Push 模式下,一个 MessageQueue 同一时刻只属于一个消费者,扩容搬运就会停摆。Pop 模式把"取消息"做成服务端的一个轻量 RPC:
- 消费者向 Proxy 发 Pop 请求,服务端从任意队列弹出一批消息,不把队列长期绑定给某个消费者。
- 消息并不是立刻删除,而是进入一个"待确认"的临时状态(checkpoint),消费者 ack 后才真正消费掉;超时没 ack 服务端自动重新可消费。
- 扩容时不再需要"队列在实例间搬家",新实例直接来 Pop 就行,没有重平衡停摆窗口。
// Pop 消费不需要再配置 MessageListenerOrderly/Concurrently 那套
// 消费者更像是一个无状态 worker,请求即消费
Consumer consumer = Provider.getConsumerBuilder()
.setConsumerGroup("trade-notify-group")
.setClientConfiguration(cfg)
.setSubscriptionExpressions(
Collections.singletonMap("trade-notify",
FilterExpression.compile("*")))
.setMessageListener(msg -> {
handle(msg);
return ConsumeResult.SUCCESS;
})
.build();
实测:同样的扩容动作
| 场景 | 4.9 Push 模式 | 5.0 Pop 模式 |
|---|---|---|
| 扩容 2 实例触发重平衡 | 积压峰值 30 万,恢复 22 分钟 | 积压峰值 1.2 万,恢复 40 秒 |
| 单实例宕机 | 该实例队列停消费约 30 秒 | 无感知,其他实例继续 Pop |
注意事项
- Pop 模式不保证顺序。我们交易通知是幂等的,无所谓;但订单状态流转这类强顺序场景还得用顺序消费,别盲目切。
- Proxy 层是一跳网络,延迟比直连 Broker 多了 1~2 ms,高频小消息要评估。
- 5.0 是 2022 年的大版本,生产落地建议先在边缘业务灰度,Proxy 的部署形态(本地/远程/集群)也要结合实际选。
小结
- 老架构的"队列绑定消费者"是扩容抖动的根因,Pop 用服务端弹出 + 无状态消费者解决了它。
- Proxy 让多语言 gRPC 客户端成为可能,对异构技术栈团队是实打实的利好。
- 顺序性要求高的业务不要为了时髦去切 Pop。
这次让我意识到,消息队列的"消费模型"比"吞吐量数字"更影响日常稳定性。Pop 不是银弹,但它确实把"扩容反而更慢"这种反直觉故障从根源上消掉了。