NATS JetStream拉取消费者并发问题:多实例避免重复消费
问题描述
我的TypeScript应用由多个项目构成,包含两类微服务:
server:REST API服务,负责向NATS发布消息,供其他服务更新MongoDBsynchronizer:专门消费NATS消息并更新数据库的进程
现在遇到的问题是,我会运行多个相同的synchronizer实例,这些实例的消费者都会绑定到同一个NATS流。我需要实现的是:这些拉取式消费者能共同消费同一个流,但绝对不能重复处理同一条消息。目前我的Consumer类实现代码如下:
import { AckPolicy, Codec, ConsumerConfig, DeliverPolicy, JetStreamClient, JetStreamManager, JsMsg, JSONCodec, ReplayPolicy, StringCodec } from "nats"; interface IEvent<T extends any> { subject: string; data: T } const createConsumerOptions = (serviceName: string): ConsumerConfig => { return { ack_policy: AckPolicy.Explicit, deliver_policy: DeliverPolicy.New, replay_policy: ReplayPolicy.Instant, ack_wait: 1000, flow_control: true, max_ack_pending: 5, max_waiting: 1, durable_name: serviceName // <- 示例:"Auth-Service" }; }; export const initConsumer = async (js: JetStreamClient, serviceName: string, subject: string) => { const jsm = await js.jetstreamManager(); if (!jsm) throw new Error("JetstreamManager Error"); await initStream(jsm, subject); const options = createConsumerOptions(serviceName); await jsm.consumers.add(subject, options).catch((ex) => console.log(ex)); let result = await js.consumers.get(subject, options); if (!result) throw new Error("Consumer Error"); return result; }; export const initStream = async (jsm: JetStreamManager, subject: string) => { const stream = await jsm.streams.find(subject).catch((ex) => console.log(ex)); if (!stream) await jsm.streams.add({ name: subject.toString(), subjects: [subject] }).catch((ex) => console.log(ex)); }; export abstract class Consumer<T extends IEvent<any>> { abstract subject: T["subject"]; abstract onMessage(data: T["data"], msg: JsMsg): Promise<any>; protected _client: Client; constructor(client: Client) { this._client = client; } async listen() { const js = await this._client.Connection?.jetstream(); if (!js) throw new Error("Jetstream Error"); const consumer = await initConsumer(js, this._client.ServiceName, this.subject); const msgs = await consumer.consume(); for await (const msg of msgs) { const data = JSONCodec().decode(msg.data) try { await this.onMessage(data, msg); msg.ack(); } catch (err) { msg.nack(); } } } }
解决方案
要实现多个synchronizer实例共享流且无重复消费,核心是利用NATS的**队列消费者(Queue Consumer)**机制,同时调整对应的配置参数:
1. 核心配置修改:添加队列名称
在createConsumerOptions函数中加入queue_name配置,相同队列名称的消费者会自动组成消费组,NATS会将消息轮询分发给组内的不同实例,从根源避免重复消费。同时调整部分参数适配多实例场景:
const createConsumerOptions = (serviceName: string): ConsumerConfig => { return { ack_policy: AckPolicy.Explicit, deliver_policy: DeliverPolicy.New, replay_policy: ReplayPolicy.Instant, ack_wait: 5000, // 适当调大ack等待时长,防止消息因处理耗时被误判为失败重发 flow_control: true, max_ack_pending: 10, // 根据实例数量和处理能力调整,避免限制消息分发 durable_name: `${serviceName}-queue`, // 保留持久化名称,确保重启后可恢复消费进度 queue_name: `${serviceName}-queue` // 关键:队列名称,实现负载均衡式消费 }; };
2. 移除限制多实例的参数
原来的max_waiting: 1会限制等待消费的客户端数量,直接导致多实例无法加入同一个消费组,必须删除该配置。
3. 其他注意事项
- Stream配置校验:确保Stream是基于目标subject创建的,多个消费者共享Stream时,Stream的消息保留策略要符合业务需求。
- 消息确认逻辑:保持
AckPolicy.Explicit不变,只有当消息处理成功时才调用msg.ack();如果处理失败,msg.nack()会让消息重新进入队列,由其他实例尝试处理。 - 消费者初始化容错:当前
initConsumer中jsm.consumers.add在队列消费者已存在时会报错,这是正常现象,保留catch逻辑忽略重复创建的错误即可。
完成上述修改后,多个synchronizer实例会自动加入同一个消费组,NATS JetStream会负责消息的均匀分发,实现同一个流的无重复消费。
内容的提问来源于stack exchange,提问作者T00rk
相关产品推荐
相关产品推荐

