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

如何用TypeScript实现NATS JetStream带通配符的Pull订阅?

JetStream通配符订阅失败问题解决

问题根源

你订阅不到消息的核心原因是流的主题配置没有覆盖到实际发布的消息主题:

  • 发布者发送的是Subject.Foo.Bar(三级主题)
  • 你创建流时用的是${this.subject}.*(假设this.subject是Subject,则流主题为Subject.*),这个通配符仅能匹配Subject.X这种二级主题,无法匹配三级的Subject.Foo.Bar,导致流根本没有存储这些消息,自然订阅不到。

解决方案

1. 修正流的主题配置

创建流时,需要使用Subject.>作为主题过滤器,>是NATS的多级通配符,能匹配所有以Subject.开头的子主题(包括任意层级的嵌套主题)。

2. 调整订阅的Pull参数

pull({ batch: 0 })的用法不推荐,batch:0表示无限制拉取,可能导致性能问题,建议设置合理的批量大小(比如batch:10)。

修正后的代码

import moment from 'moment';
import { consumerOpts, JetStreamPullSubscription, JsMsg, Msg, StringCodec } from 'nats';

import { IEvent } from '../events/event';
import { NatsClient } from './nats-client';

export abstract class NatsNestedListener<S extends IEvent>
{
  abstract onMessage(data: S['data'], msg: JsMsg): Promise<void>;
  abstract subject: S['subject'];  
  
  private client: NatsClient;
  
  protected subscription: JetStreamPullSubscription;

  constructor(client: NatsClient)
  {
    this.client = client;
  }

  private subscriptionOptions()
  {
    const options = consumerOpts();

    options.ackWait(5000);
    options.deliverAll();
    options.startSequence(1);
    options.durable(this.client.QueueGroupName);
    options.manualAck();

    return options;
  }

  async listen()
  {
    // 修正:用多级通配符创建流,确保覆盖所有嵌套主题
    await this.client.initStream(`${this.subject}.>`);

    // 订阅时可以用Subject.*或Subject.>,根据需求选择
    this.subscription = await this.client.JetStreamClient.pullSubscribe(`${this.subject}.>`, this.subscriptionOptions());
    // 修正:设置合理的批量大小,比如10
    this.subscription.pull({ batch: 10 });
  
    for await(const jsMessage of this.subscription)
    {
      console.log('Message Received', `[Subject]: ${jsMessage.subject} - [Seq]: ${jsMessage.seq}`);

      await this.processMessage(jsMessage); // 添加await,确保消息处理完成再拉取下一批
      this.subscription.pull({ batch: 10 });
    }   
  }

  parseMessage(msg: JsMsg | Msg)
  {
    const stringCodec = StringCodec();
    const decodedData = stringCodec.decode(msg.data)

    return JSON.parse(decodedData);
  }

  async processMessage(msg: JsMsg)
  {
    const data = this.parseMessage(msg);
    await this.onMessage(data, msg); // 添加await,确保业务逻辑执行完成
  }
}

关键注意事项

  • JetStream流的主题匹配逻辑:流的主题过滤器决定了哪些消息会被存储,必须确保它覆盖你要订阅的所有消息主题。*匹配单一级别,>匹配所有后续级别。
  • Pull订阅的批量设置:避免使用batch:0,建议根据业务场景设置合适的批量值,平衡性能和消息处理速度。
  • 异步处理的await:在调用异步的processMessage和onMessage时添加await,确保消息处理完成后再进行下一次拉取或确认,避免消息丢失或重复处理。

内容的提问来源于stack exchange,提问作者T00rk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 09:01:07