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

为Azure Service Bus启用会话的订阅创建持续接收器

解决Azure Service Bus启用会话主题的持续监听问题

问题场景

我们有一个启用会话的Azure Service Bus主题订阅,消息带有动态多样的会话ID,无法提前指定会话ID使用acceptSession方法。使用createReceiver会抛出运行时错误,而acceptNextSession仅能处理第一个会话的消息,无法持续获取后续会话的消息。当前使用的是新版本的@azure/service-bus npm包,旧示例中的方法已不可用。

解决方案

在新版本@azure/service-bus(7.x及以上)中,要持续处理所有动态会话的消息,需要通过**循环调用acceptNextSession**来不断获取下一个可用的会话,并为每个会话创建独立的接收器进行消息处理。具体实现如下:

import { ServiceBusClient, ServiceBusReceiver } from "@azure/service-bus";

async function startSessionListener(connectionString: string, topicName: string, subscriptionName: string, handleTopicMessage: (message: any) => Promise<void>) {
  const sbClient = new ServiceBusClient(connectionString);
  const logger = console; // 根据实际情况替换为你的日志实例

  // 单个会话的消息处理逻辑
  async function processSession(receiver: ServiceBusReceiver) {
    try {
      await receiver.subscribe({
        processMessage: async (message) => {
          await handleTopicMessage(message);
          // 根据业务需求完成消息,避免重复消费
          await receiver.completeMessage(message);
        },
        processError: async (error) => {
          logger.error(`会话 ${receiver.sessionId} 处理错误: ${error.errorMessage}`);
          // 针对消息级错误做死信处理
          if (error.errorSource === "message") {
            await receiver.deadLetterMessage(error.message!, { reason: "消息处理失败" });
          }
        }
      });

      // 等待会话自然结束(无新消息或会话被关闭)
      await receiver.getSessionState();
      await receiver.close();
      logger.info(`会话 ${receiver.sessionId} 处理完成并关闭`);
    } catch (err) {
      logger.error(`会话 ${receiver.sessionId} 处理异常: ${err}`);
      // 确保异常时关闭接收器释放资源
      await receiver.close().catch(() => {});
    }
  }

  // 循环获取并处理下一个会话
  async function listenForSessions() {
    while (true) {
      try {
        const receiver = await sbClient.acceptNextSession(topicName, subscriptionName);
        logger.info(`获取到新会话: ${receiver.sessionId}`);
        // 异步处理当前会话,不阻塞后续会话的获取
        processSession(receiver);
      } catch (err) {
        logger.error(`获取会话失败: ${err}`);
        // 临时错误后延迟重试,避免高频报错
        await new Promise(resolve => setTimeout(resolve, 5000));
      }
    }
  }

  // 启动持续监听
  await listenForSessions();
}

// 使用示例
const serviceBusSettings = { connectionString: "你的Service Bus连接字符串" };
const topicName = "目标主题名称";
const subscriptionName = "目标订阅名称";

async function handleTopicMessage(message: any) {
  console.log(`处理消息: ${JSON.stringify(message.body)}, 会话ID: ${message.sessionId}`);
}

startSessionListener(serviceBusSettings.connectionString, topicName, subscriptionName, handleTopicMessage)
  .catch(err => console.error(`监听启动失败: ${err}`));

关键要点说明

  • 循环获取会话:通过while(true)循环持续调用acceptNextSession,每次获取到可用会话后异步启动处理逻辑,不阻塞后续会话的获取。
  • 会话独立隔离:每个会话对应独立的接收器,处理完成或出错后主动关闭,确保资源正常释放。
  • 消息状态管理:根据业务场景调用completeMessage、abandonMessage或deadLetterMessage,避免消息重复消费或丢失。
  • 异常容错:在获取会话失败时添加延迟重试逻辑,防止临时网络或服务错误导致监听中断。

原代码无效原因

之前的代码仅调用了一次acceptNextSession,只能处理第一个被获取到的会话。当该会话的消息处理完成后,没有逻辑触发后续会话的获取,因此无法持续接收新会话的消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 08:40:32