Blog

普通Pub/Sub、Queue Group和Request-Reply浅析

1435597771 · Aug 4, 2026

普通Pub/Sub订阅模型 普通 Pub/Sub 是 Core NATS 最基础的通信模型。Publisher 向 Subject 发布消息,所有匹配且当前活动的独立 Subscriber 都可以获得一份消息。 但是普通的Pub/Sub订阅模型无法解决下面两个场景 多实例场景:邮件服务有两个实例, 都订阅jobs.email 。发布一条发送邮件任务后,两个实例都会收到消息,用户可能会收到两封邮件 等待响应状态:订单服务查询订单状态,他需要在发布消息后,等待订单状态的返回。 为了解决多实例负载分担,Core NATS 提供了 Queue Group;为了让请求方等待明确响应,NATS 在 Pub/Sub 之上提供了 Request-Reply 通信模式。 Queue Group:组内分担,组间广播 从消息分发的逻辑视角看,NATS Server 会把一个 Queue Group 当成一个整体;但 Queue Group 本身不会处理消息,最终仍由组内某一个具体 Subscriber 接收并执行回调。当消息发布时,NATS Server 会为每个匹配的 Queue Group 选择一个活动成员,并将该消息投递给这个成员。 QueueGroup的核心代码是: _, err := nc.QueueSubscribe(   "jobs.email",   "email-workers",   handler, ) 我们需要说明这三个要素: jobs.email : Subject,消息路由,决定订阅什么消息 email-workers: Queue Group的名字,只要订阅了相同的Subject,并且组名也相同,那么就是在一个组内 handler:回调函数 组内分担是如何分担的: 判定一个组的规则是相同的Subject+相同的QueueName那么就会被视为同一个组,这个组可以看作是一个逻辑组。因为当NATS在投递消息的时候,他是可以看到这个组内所有的worker,然后选择其中一个投递消息。 他的特性是: 动态扩缩容:可以随时增加或者减少成员 随机选择:从应用设计角度,不应依赖具体选择算法。Queue Group 只保证每条消息在同组内交给一个成员,不保证严格轮询、严格随机或最终均分。 至多一次:如果消息已经投递给某个 Worker,而 Worker 在业务处理中崩溃,Core NATS 不会确认数据库操作是否成功,也不保证把这条旧消息重新投递给其他成员。 非持久化:如果整个 Queue Group 在消息发布时都没有活动成员,Core NATS 不会因为组名存在而保存消息;成员之后重新上线,也不会补收离线期间的消息。 组间广播的含义 我们说了要是同一个组,那么必须是相同的subject+相同的组名,那么如果是相同的subject不同的组名呢。 从组与组之间的关系看,不同 Queue Group 会像彼此独立的订阅者一样,各自获得一份消息;但在每个组内部,这一份消息仍然只投递给一个成员。这也就意味着,实际上在不同的组,其实就是不同的subscriber订阅了相同的subject而已。可以给组间广播看作是普通的pub/sub模型。 我们了解组内分担和组间广播是如何实现的,那么就知道了Queue Group解决多实例的方法,可以看作是 一个组收到消息,然后选择组内一个成员处理消息。 Request-Reply是什么 我们说为了满足等待响应状态这个功能,需要用到Request-Reply这个订阅模型,那么Request-Reply长什么样子呢 在客户端: msg, err := nc.Request("orders.inventory.check",   []byte(`{"order_id":"ord_8w2k"}`),   2*time.Second) ​ if err != nil {   // 处理超时或无响应者   return } fmt.Println("收到回复:", string(msg.Data)) 在服务端: _, err := nc.QueueSubscribe(   "orders.inventory.check",   "inventory-services",   func(msg *nats.Msg) {       reply := []byte(`{"in_stock":true,"warehouse":"us-east"}`) ​       if err := msg.Respond(reply); err != nil {           log.Printf("respond failed: %v", err)       }   }, ) 在使用上我们看到过程是很简单的: 客户端使用Request方法,向orders.inventory.check这个subject中发送消息,并设置等待时间2s 服务端收到消息,处理业务 服务端用msg.Respond 向客户端发送消息 客户端收到回复或者超时 整个过程很简单,但是我们需要知道服务端是如何向客户端发送请求的,我们需要给Request这个方法拆开。 Request使用了Reply这个字段,在Request内部,他会先创建一个inbox,用来当作服务端响应的subject,然后客户端会先订阅这个inbox,所以在客户端发送消息之前,他本身就已经订阅了inbox这个主题。 Request的伪代码可以看作是: function Request(subject, payload, timeout):   // ========== 第一步:生成唯一的 Inbox ==========   // 示例:_INBOX.ABC123.xyz789   inbox = "_INBOX." + connection_id + "." + random_unique_token()       // ========== 第二步:订阅这个 Inbox ==========   subscription = Subscribe(inbox) ​   // ========== 第三步:发布请求,并附带 Inbox 地址 ==========   message = {       subject: subject,         // 目标主题,比如 "orders.inventory.check"       reply:   inbox,           // 关键!把 Inbox 写进 reply 字段       data:   payload   }   Publish(message) ​   // ========== 第四步:等待回复 ==========   try:       reply_msg = WaitForMessage(subscription, timeout)       return reply_msg.data     // 成功拿到响应   except TimeoutError:       raise "请求超时"   except NoRespondersError:       raise "没有响应者" // NATS Server 处理本次请求时,没有发现请求 Subject 的可用订阅者,因此可能快速返回 No Responders,而不是等待完整超时时间。   finally:       // ========== 第五步:清理 ==========       Unsubscribe(subscription) // 释放 Inbox 订阅 那么Respond就比较好理解了,Respond就是向Inbox这个Subject发送消息就好了。 所以Request-Reply可以看作是两个普通的Pub/Sub消息订阅模型。 在不同的阶段客户端和服务端的角色会发生变化 请求阶段:客户端是 Publisher,服务端是 Subscriber。 响应阶段:服务端是 Publisher,客户端是 Subscriber。 以上就是对普通Pub/Sub,Queue Group和Request-Reply的简单介绍,希望对你有所帮助。