You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

NATS JetStream拉取消费者并发问题:多实例避免重复消费

问题描述

我的TypeScript应用由多个项目构成,包含两类微服务:

  • server:REST API服务,负责向NATS发布消息,供其他服务更新MongoDB
  • synchronizer:专门消费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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 04:59:53