NATS queue groups:每条消息随机挑一个成员,组名打错会静默多出一个组
信源:NATS 官方文档 · Queue Groups 概念、Core NATS Queue Groups、NATS CLI(docs.nats.io,2.15)。
它解决什么
标准 pub/sub 里每个订阅者都收到每条消息。仓储(warehouse)场景不行:一个订单只能被一个打包员处理,不能零次也不能两次。三个 warehouse 副本跑普通订阅,三个都会打包同一个订单。
queue group 的定义:同一 subject 上、共享一个名字的一组订阅者。服务器把整组当成一个逻辑订阅者,每条消息只在组里挑恰好一个成员投递。
组名由应用自己定,订阅时带上;服务器不做配置。第一个成员订阅时服务器就知道这个组存在。
CLI 里把普通订阅改成队列组成员:
nats sub orders.created --queue packers
--queue 就是把普通订阅变成队列组成员。多个终端跑同一条命令就是多个成员。
服务器怎么挑成员
服务器把组内活跃成员放在一个列表里,每条消息随机取一个下标投递。单机上对可用成员是均匀随机;集群会加 locality 偏好。
随机的后果:服务器不像轮询那样公平轮转,同一个成员可能连续被选中,一小批消息看分布会歪,量大了才平。
成员动态增减
订阅即加入,取消订阅/断连即离开,无需配置或协调。加油流量时加成员立刻分担;停掉一个,服务器把它移出列表,下一条给剩下的。适合自动扩缩容:池子从 10 缩到 2,broker 不用改配置。
一个硬限制:核心 NATS 是 at-most-once。服务器选了一个成员、投递之后它挂了,这条消息就没了,服务器不会转投给别的成员。要扛 worker 崩溃的活儿要放进 JetStream 的持久工作队列。
队列成员和普通订阅者共存
同一 subject 上,queue group 和普通订阅者互不干扰,两套分发独立跑。analytics 用普通订阅收到每一条;packers 组里恰好一个成员收到。同一 subject 同时承载两种行为。
一个组只在一个 subject 内分摊
成员匹配发生在 subject 匹配之后。服务器先找出 subject 匹配的订阅,再在其中做组挑选。同名但不同 subject 的成员不分摊负载:packers 组里一个订 orders.created、一个订 orders.shipped,消息只会落到匹配的那一个。名字跨 subject 不起作用。
组可以用通配符订阅,然后在该通配符匹配到的所有 subject 上分摊(如 orders.*.created)。
跨区域放置
同一个组在多个集群有成员时,服务器优先选发布者所在集群的成员,避免跨区流量。这是多集群(super-cluster)下 queue group 的 geo-affinity,本机单服务器场景用不到。
坑
- 组名打错会静默多出一个组。 服务器按精确字符串匹配,
packers和packer是同一 subject 上的两个组,两个订阅都不报错,每条消息各投给每组一个成员——工作被重复处理,而不是负载均衡。要每个成员逐字节一致。 - 别指望有序或均匀切分。选择是按消息随机的,不是轮询,短时突发会歪。需要按客户严格有序的活儿,放单个订阅者,别用 queue group。
- 让成员的工作可重复。成员慢或被短暂切断时,仍可能在处理发布者以为已丢的消息,而核心 NATS 不会重发。按
order_id打包,已处理的跳过。
与 request-reply 结合
queue groups 让服务横向扩展不需要 service mesh 或 API 网关:每个请求恰好到达一个实例。服务代码不用知道其它实例、不用 leader 选举、不用协调。
nats reply api.calculate --queue api-workers "..."
请求方用 nats request,负载自动分摊到各实例。
动态扩缩容示例
先起一个成员:
nats sub tasks --queue workers
负载涨了在新终端再加;负载降了 Ctrl+C 停掉,剩下的自动接手。
代码侧就是每实例启动时 queueSubscribe,退出时 unsubscribe。Python 示例里在发布前用 await nc.flush(timeout=1) 等服务器确认所有订阅到位,避免第一条消息在订阅生效前发出而丢失。
与 Kafka / RabbitMQ 的对照(官方 glossary)
- 配置:NATS queue group 无需服务端配置;传统队列常需预配置。
- 扩展:NATS 随订阅者数量动态调整;Kafka 需分区/手动调优,RabbitMQ 需要队列/消费者配置。
- 持久化:NATS 默认无状态,可用 JetStream 加持久;传统队列默认持久、开销更大。
- 延迟:NATS 实时低延迟;持久化系统因持久化和协调更高。
对比消费者组:Kafka 需要显式分区和 offset 管理来达到类似目的,NATS queue group 复杂度更低。
站内多数消息队列选型文比较的是持久化、吞吐和运维成本,这里更值得记住的是另一个维度:NATS queue group 的投递语义、成员选择方式和失败边界都写在服务器行为里,而不是写在配置里——组名、subject 范围、at-most-once 这三件事决定了它能干什么、不能干什么。