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

Azure Service Bus:如何像命名ServiceBusClients那样注入多ServiceBusReceiver单例

解决Azure Service Bus多Receiver/Sender的DI注册问题

不用为每个ServiceBusReceiver或ServiceBusSender单独写包装类,Azure SDK已经提供了简洁的扩展方法来处理多实例注册,以下是几种实用方案:

方案1:直接注册带标识的单例实例(.NET 8+推荐)

利用.NET 8新增的键注入功能,为每个队列/主题的Receiver/Sender指定唯一标识,注册时关联到根ServiceBusClient:

services.AddAzureClients(builder =>
{
    // 注册全局单例ServiceBusClient
    builder.AddServiceBusClient(myConnectionString);

    // 为每个队列注册带键的ServiceBusReceiver单例
    builder.AddServiceBusReceiver("queue-a").WithKey("QueueA");
    builder.AddServiceBusReceiver("queue-b").WithKey("QueueB");
    builder.AddServiceBusReceiver("queue-c").WithKey("QueueC");

    // Sender实例同理
    builder.AddServiceBusSender("topic-order").WithKey("OrderTopicSender");
    builder.AddServiceBusSender("queue-notify").WithKey("NotifyQueueSender");
});

然后在BackgroundService中通过[FromKeyedServices]特性注入对应实例:

public class QueueABackgroundService : BackgroundService
{
    private readonly ServiceBusReceiver _queueAReceiver;

    // 通过键标识注入指定的Receiver
    public QueueABackgroundService([FromKeyedServices("QueueA")] ServiceBusReceiver receiver)
    {
        _queueAReceiver = receiver;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // 监听队列消息
        await foreach (var message in _queueAReceiver.ReceiveMessagesAsync(stoppingToken))
        {
            try
            {
                // 处理消息逻辑
                await _queueAReceiver.CompleteMessageAsync(message, stoppingToken);
            }
            catch (Exception ex)
            {
                // 异常处理,比如放弃或死信消息
                await _queueAReceiver.DeadLetterMessageAsync(message, stoppingToken);
            }
        }
    }
}

方案2:使用命名客户端兼容旧版本

如果你的项目基于.NET 6/7,可以通过命名ServiceBusClient来关联对应的Receiver/Sender,同样实现单例注册:

services.AddAzureClients(builder =>
{
    // 为每个队列/主题注册命名的ServiceBusClient
    builder.AddServiceBusClient(myConnectionString).WithName("QueueA_Client");
    // 基于命名客户端创建Receiver
    builder.AddServiceBusReceiver("QueueA_Client", "queue-a");

    builder.AddServiceBusClient(myConnectionString).WithName("QueueB_Client");
    builder.AddServiceBusReceiver("QueueB_Client", "queue-b");
});

注入时通过IAzureClientFactory<ServiceBusReceiver>获取指定命名的实例:

public class QueueBBackgroundService : BackgroundService
{
    private readonly ServiceBusReceiver _queueBReceiver;

    public QueueBBackgroundService(IAzureClientFactory<ServiceBusReceiver> receiverFactory)
    {
        _queueBReceiver = receiverFactory.CreateClient("QueueB_Client");
    }

    // 消息监听逻辑...
}

方案3:自定义工厂类(灵活扩展场景)

如果需要动态创建Receiver/Sender(比如根据配置加载队列列表),可以写一个简单的工厂类来管理实例,确保单例复用:

public class ServiceBusEntityFactory
{
    private readonly ServiceBusClient _rootClient;
    private readonly Dictionary<string, ServiceBusReceiver> _receiverCache = new();
    private readonly Dictionary<string, ServiceBusSender> _senderCache = new();
    private readonly object _lockObj = new();

    public ServiceBusEntityFactory(ServiceBusClient client)
    {
        _rootClient = client;
    }

    public ServiceBusReceiver GetReceiver(string queueName)
    {
        lock (_lockObj)
        {
            if (!_receiverCache.TryGetValue(queueName, out var receiver))
            {
                receiver = _rootClient.CreateReceiver(queueName);
                _receiverCache[queueName] = receiver;
            }
            return receiver;
        }
    }

    public ServiceBusSender GetSender(string entityName)
    {
        lock (_lockObj)
        {
            if (!_senderCache.TryGetValue(entityName, out var sender))
            {
                sender = _rootClient.CreateSender(entityName);
                _senderCache[entityName] = sender;
            }
            return sender;
        }
    }
}

注册工厂为单例:

services.AddSingleton<ServiceBusEntityFactory>();

在BackgroundService中使用:

public class QueueCBackgroundService : BackgroundService
{
    private readonly ServiceBusReceiver _queueCReceiver;

    public QueueCBackgroundService(ServiceBusEntityFactory factory)
    {
        _queueCReceiver = factory.GetReceiver("queue-c");
    }

    // 消息处理逻辑...
}

核心注意事项

  • 所有ServiceBusReceiver和ServiceBusSender必须保持单例生命周期,避免频繁创建实例导致的连接资源泄漏。
  • 不要手动调用Dispose方法,这些实例由DI容器管理,应用关闭时会自动清理。
  • 当注册多个同类型实例时,必须通过键或命名客户端区分,否则DI容器无法解析正确的实例。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 15:15:04