共享队列组(Shared Queue Group)
这是 RobustMQ AMQP 实现里最核心的设计:每个 AMQP 队列在内部就是一个共享消费组,而共享消费组这套基础设施是和 MQTT 共享订阅、NATS 队列组共用的——AMQP 没有另起一套队列引擎。
为什么复用共享消费组
AMQP 的"一个队列、多个消费者竞争消费"和 MQTT/NATS 里"一组客户端共享订阅同一个 topic"本质上是同一个问题:一份数据,多个消费者中只有一个该拿到每一条。RobustMQ 已经为 MQTT/NATS 实现了这套"选一个 leader 节点负责拉取和分发"的机制,AMQP 队列直接复用它,而不是重新发明。
队列如何映射为共享消费组
- 队列名即共享消费组的组名(
group_name = queue_name)。 - 组的创建是按需触发的:队列被
Queue.Declare时只创建元数据和存储分片,真正的共享消费组在第一次有客户端Basic.Get或Basic.Consume这个队列时才创建。 - 组的 leader 由 meta-service 按集群负载选举,选举结果(以及创建组这次操作的结果)持久化在 Raft,对所有节点可见。
Leader 解析与创建时的一致性
第一次访问某个队列时,RobustMQ 需要拿到(或创建)它的共享消费组以确定 leader。这里有一个容易踩的坑:如果创建完组之后再单独发一次"读"去确认刚创建的内容,读操作可能会落在一个还没同步到最新数据的节点上(Raft 的"提交"和"应用到本地状态机"是两个阶段,不同节点的应用进度可能有短暂差异)。
RobustMQ 的做法是让"创建"操作本身直接返回权威结果——创建共享消费组的 meta-service 请求在处理时已经算出了完整信息(包括选出的 leader),直接把这份数据带回给调用方,调用方不需要再额外发一次读请求去确认。这避免了"写完立刻读、读到还没应用的旧状态"的竞态。
竞争消费(Competing Consumers)
多个消费者(可能连在不同节点上)对同一个队列 Basic.Consume,就是这个共享消费组的多个成员:
- 队列 leader 节点上的推送任务按 round-robin 在组内成员之间轮询,把下一条待投递消息交给下一个"当前有空闲 prefetch 配额"的成员。
- 如果被选中的成员因为暂时没有配额(prefetch 已满)或连接不可用而无法接收,推送任务会尝试组内下一个成员,而不是阻塞等待。
- 如果消费者在推送任务所在的节点(leader)上,直接本地投递;否则通过一次 gRPC(
SendShareGroupMessage)把消息发到消费者实际所在的节点和连接。
Basic.Get 走的是同一套 leader 机制:请求先解析出队列当前的 leader,若不是本节点则转发(FetchAmqpQueueMessage),对客户端完全透明。
已知限制:跨节点 QoS
Basic.Qos 的 prefetch 限制,是由队列 leader 节点检查消费者当前的未确认消息数来决定是否可以继续投递的。这个检查只在消费者与队列 leader 同一节点时是强一致的——因为未确认计数就存在那个节点的本地状态里。如果消费者连接在别的节点上,leader 节点看不到它实时的未确认计数,prefetch 限制目前是尽力而为,不做跨节点强制同步(这需要引入额外的 RPC 查询或状态复制,目前尚未实现)。
