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

如何通过代码创建Azure Service Bus topic订阅并查询现有订阅

基于连接字符串操作Azure Service Bus Topic订阅实现方案

不需要目标Service Bus namespace的Azure Portal访问权限,只要对方提供的连接字符串对应SAS策略开启Manage权限,即可直接通过官方SDK完成订阅列举、新建、删除操作,无需调用ARM接口,适配Console App、Function App等任意运行形态。以下为Java、JavaScript两个技术栈的实现。

前置校验:拿到连接字符串后先确认权限级别,仅带Send/Listen权限的连接字符串只能做消息收发,无法执行订阅实体的管理操作,会直接返回401错误。

JavaScript/Node.js 实现

依赖安装

执行命令安装官方最新版SDK:
npm install @azure/service-bus
不要使用已废弃的旧版azure-sb包,该版本存在兼容性问题且不再维护

核心代码实现

const { ServiceBusAdministrationClient } = require("@azure/service-bus");
// 替换为对方提供的Service Bus连接字符串
const connectionString = "YOUR_SERVICE_BUS_CONNECTION_STRING";
const adminClient = new ServiceBusAdministrationClient(connectionString);

/**
 * 列举指定Topic下的全量订阅
 * @param {string} topicName 目标Topic名称
 */
async function listAllSubscriptions(topicName) {
  const subscriptions = [];
  const subscriptionIterator = adminClient.listSubscriptions(topicName);
  for await (const sub of subscriptionIterator) {
    subscriptions.push({
      name: sub.subscriptionName,
      createdAt: sub.createdAt,
      lastAccessedAt: sub.accessedAt
    });
  }
  return subscriptions;
}

/**
 * 新建Topic订阅
 * @param {string} topicName 目标Topic名称
 * @param {string} subscriptionName 待创建的订阅名称
 * @param {object} config 订阅配置,可选
 */
async function createSubscription(topicName, subscriptionName, config = {}) {
  return adminClient.createSubscription(topicName, subscriptionName, {
    maxDeliveryCount: config.maxDeliveryCount || 10,
    defaultMessageTimeToLive: config.messageTtl || "P7D", // ISO 8601时长格式,默认7天
    ...config
  });
}

/**
 * 批量删除过期订阅
 * @param {string} topicName 目标Topic名称
 * @param {number} expireDayThreshold 超过指定天数未访问即判定为过期
 */
async function deleteExpiredSubscriptions(topicName, expireDayThreshold = 30) {
  const allSubscriptions = await listAllSubscriptions(topicName);
  const currentTimestamp = Date.now();
  for (const sub of allSubscriptions) {
    const lastAccessTimestamp = new Date(sub.lastAccessedAt).getTime();
    const idleDays = (currentTimestamp - lastAccessTimestamp) / (1000 * 60 * 60 * 24);
    if (idleDays > expireDayThreshold) {
      await adminClient.deleteSubscription(topicName, sub.name);
      console.log(`已删除过期订阅:${sub.name},最后访问时间:${sub.lastAccessedAt}`);
    }
  }
}

// 调用示例
async function main() {
  const topicName = "YOUR_TOPIC_NAME";
  // 列举所有订阅
  const subs = await listAllSubscriptions(topicName);
  console.log("当前订阅列表:", subs);
  // 新建订阅
  await createSubscription(topicName, "test-sub-001");
  // 删除30天未访问的过期订阅
  await deleteExpiredSubscriptions(topicName, 30);
}

main().catch(console.error);

该实现可直接运行在本地Console环境,部署到Azure Function App时只需将连接字符串、Topic名称配置为应用配置项,不需要绑定任何Azure账号权限,只要运行网络能连通对方Service Bus端点即可。

Java 实现

依赖引入

在pom.xml中引入官方SDK依赖:

<dependency>
    <groupId>com.azure</groupId>
    <artifactId>azure-messaging-servicebus</artifactId>
    <version>7.14.0</version>
</dependency>

可自行替换为最新正式版本,避免使用旧版azure-servicebus包

核心代码实现

import com.azure.messaging.servicebus.administration.ServiceBusAdministrationClient;
import com.azure.messaging.servicebus.administration.ServiceBusAdministrationClientBuilder;
import com.azure.messaging.servicebus.administration.models.CreateSubscriptionOptions;
import com.azure.messaging.servicebus.administration.models.SubscriptionProperties;
import java.time.Duration;
import java.time.OffsetDateTime;
import java.util.ArrayList;
import java.util.List;

public class ServiceBusSubscriptionManager {
    private static final String CONNECTION_STRING = "YOUR_SERVICE_BUS_CONNECTION_STRING";
    private final ServiceBusAdministrationClient adminClient;

    public ServiceBusSubscriptionManager() {
        this.adminClient = new ServiceBusAdministrationClientBuilder()
                .connectionString(CONNECTION_STRING)
                .buildClient();
    }

    /**
     * 列举指定Topic下全量订阅
     * @param topicName 目标Topic名称
     * @return 订阅属性列表
     */
    public List<SubscriptionProperties> listAllSubscriptions(String topicName) {
        List<SubscriptionProperties> subscriptionList = new ArrayList<>();
        adminClient.listSubscriptions(topicName).forEach(subscriptionList::add);
        return subscriptionList;
    }

    /**
     * 新建Topic订阅
     * @param topicName 目标Topic名称
     * @param subscriptionName 待创建的订阅名称
     * @param maxDeliveryCount 最大投递次数
     * @param messageTtl 消息默认存活时间
     */
    public void createSubscription(String topicName, String subscriptionName, int maxDeliveryCount, Duration messageTtl) {
        CreateSubscriptionOptions options = new CreateSubscriptionOptions()
                .setMaxDeliveryCount(maxDeliveryCount)
                .setDefaultMessageTimeToLive(messageTtl);
        adminClient.createSubscription(topicName, subscriptionName, options);
    }

    /**
     * 批量删除过期订阅
     * @param topicName 目标Topic名称
     * @param expireDayThreshold 超过指定天数未访问即判定为过期
     */
    public void deleteExpiredSubscriptions(String topicName, int expireDayThreshold) {
        List<SubscriptionProperties> allSubscriptions = listAllSubscriptions(topicName);
        OffsetDateTime currentTime = OffsetDateTime.now();
        for (SubscriptionProperties sub : allSubscriptions) {
            Duration idleDuration = Duration.between(sub.getAccessedAt(), currentTime);
            if (idleDuration.toDays() > expireDayThreshold) {
                adminClient.deleteSubscription(topicName, sub.getSubscriptionName());
                System.out.printf("已删除过期订阅:%s,最后访问时间:%s%n", sub.getSubscriptionName(), sub.getAccessedAt());
            }
        }
    }

    // 调用示例
    public static void main(String[] args) {
        ServiceBusSubscriptionManager manager = new ServiceBusSubscriptionManager();
        String topicName = "YOUR_TOPIC_NAME";
        // 列举所有订阅
        List<SubscriptionProperties> subs = manager.listAllSubscriptions(topicName);
        subs.forEach(sub -> System.out.println("订阅名称:" + sub.getSubscriptionName()));
        // 新建订阅
        manager.createSubscription(topicName, "test-sub-001", 10, Duration.ofDays(7));
        // 删除30天未访问的过期订阅
        manager.deleteExpiredSubscriptions(topicName, 30);
    }
}

补充说明:

  • 上述代码已内置分页处理逻辑,列举订阅时会自动拉取全量数据,不需要手动处理分页参数
  • 如果需要为订阅配置消息过滤规则(SqlFilter、CorrelationFilter),两个SDK均提供对应接口,可在创建订阅时传入规则配置
  • 不管部署在本地Console、虚拟机还是Function App,都不需要持有对方Azure订阅的RBAC权限,仅靠连接字符串即可完成鉴权
  • 如果对方Service Bus开启了IP白名单,需要将你方程序运行的出口IP加入白名单,否则会出现网络连接超时

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 06:36:37