普通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的简单介绍,希望对你有所帮助。