普通Pub/Sub、Queue Group和Request-Reply浅析
1435597771 ·
普通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的简单介绍,希望对你有所帮助。