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

.NET 8中Singleton服务结合Hosted Services多次实例化问题排查

.NET 8中Azure Service Bus消费者注册框架的Singleton实例问题

问题描述

我正在.NET 8中构建类似Masstransit的Azure Service Bus消费者注册框架,遇到以下问题:

  • 通过ServiceBusManager.AddConsumers注册消费者,将处理器加入_processors列表;
  • 运行Hosted Service时,获取同一ServiceBusManager实例调用StartProcessingAsync,但此时_processors列表为空。

尽管已将ServiceBusManager注册为Singleton,它却被多次实例化,怀疑这与服务解析、配置方式(尤其是和Hosted Service或连接构建器的交互)有关。

注册代码示例

builder.Services.AddSingleton<IServiceBusManager, ServiceBusManager>();
builder.Services.AddConsumerServiceBusConnection(x =>
    {
        var topicName = builder.Configuration["ServiceBusSettings:TopicName"];
        var subscriptionName = builder.Configuration["ServiceBusSettings:SubscriptionName"];
         x.AddConsumer<ShipmentCreatedConsumer>(topicName, subscriptionName);
         x.AddConsumer<TestConsumer>("testtopic", "testsuscription");
    });

扩展方法实现

public static class ServiceCollectionExtensions
{
    public static IServiceCollection AddConsumerServiceBusConnection(
        this IServiceCollection services,
        Action<ServiceBusConnectionBuilder> configure)
    {
        // Build the service provider to resolve the singleton instance
        using (var serviceProvider = services.BuildServiceProvider())
        {
            var serviceBusManager = serviceProvider.GetRequiredService<IServiceBusManager>();
            var builder = new ServiceBusConnectionBuilder(serviceBusManager);
            configure(builder);
        }
        return services;
    }
}

ServiceBusConnectionBuilder实现

public class ServiceBusConnectionBuilder
{
    private readonly IServiceBusManager _serviceBusManager;

    public ServiceBusConnectionBuilder(IServiceBusManager serviceBusManager)
    {
        _serviceBusManager = serviceBusManager;
    }

    public ServiceBusConnectionBuilder AddConsumer<TConsumer>(string topicName, string suscrptionName) 
        where TConsumer : IServiceBusConsumer
    {
        _serviceBusManager.AddConsumer<TConsumer>(topicName, suscrptionName);
        return this;
    }
}

ServiceBusManager实现

public class ServiceBusManager : IServiceBusManager
{
    private readonly ServiceBusClient _serviceBusClient;
    private readonly IServiceProvider _serviceProvider;
    private readonly ConcurrentDictionary<string, ServiceBusProcessor> _processors;

    public ServiceBusManager(ServiceBusClient serviceBusClient,
        IServiceProvider serviceProvider,
        IOptions<ServiceBusSettings> serviceBusSettings)
    {
        _serviceBusClient = serviceBusClient ?? throw new ArgumentNullException(nameof(serviceBusClient));
        _serviceProvider = serviceProvider ?? throw new ArgumentNullException(nameof(serviceProvider));
        _processors = new ConcurrentDictionary<string, ServiceBusProcessor>();
    }

    public void AddConsumer<TConsumer>(string topicName, string subscriptionName) where TConsumer : IServiceBusConsumer
    {
        if (string.IsNullOrWhiteSpace(topicName))
            throw new ArgumentException("Topic name cannot be null or empty.", nameof(topicName));

        if (string.IsNullOrWhiteSpace(subscriptionName))
            throw new ArgumentException("Subscription name cannot be null or empty.", nameof(subscriptionName));

        var processor = _serviceBusClient.CreateProcessor(topicName, subscriptionName, new ServiceBusProcessorOptions
        {
            AutoCompleteMessages = false,
            MaxConcurrentCalls = 1,       
            PrefetchCount = 10            
        });

        processor.ProcessMessageAsync += async args =>
        {
            using var scope = _serviceProvider.CreateAsyncScope();
            var consumer = scope.ServiceProvider.GetRequiredService<TConsumer>();
            await consumer.ProcessMessage(args);
        };

        processor.ProcessErrorAsync += args =>
        {
            using var scope = _serviceProvider.CreateAsyncScope();
            var consumer = scope.ServiceProvider.GetRequiredService<TConsumer>();
            return consumer.ProcessError(args);
        };

        if (!_processors.TryAdd($"{topicName}:{subscriptionName}", processor))
        {
            throw new InvalidOperationException($"Consumer for {topicName}:{subscriptionName} is already registered.");
        }
    }

    public async Task StartProcessingAsync(CancellationToken cancellationToken)
    {
        var startTasks = _processors.Values.Select(processor => processor.StartProcessingAsync(cancellationToken));
        await Task.WhenAll(startTasks);
    }

    public async Task StopProcessingAsync()
    {
        var stopTasks = _processors.Values.Select(processor => processor.StopProcessingAsync());
        await Task.WhenAll(stopTasks);
    }
}

消费者接口定义

public interface IServiceBusConsumer
{
    Task ProcessMessage(ProcessMessageEventArgs args);
    Task ProcessError(ProcessErrorEventArgs args);
}

Hosted Service实现

public class ServiceBusHostedService : CronJobServiceBase
{
    private readonly IServiceProvider _serviceProvider;
    private AsyncServiceScope _scope;
    private IHostedServiceTask _taskService;

    public ServiceBusHostedService(
        IOptions<ServiceBusHostedServiceSettings> hostedServiceSettings
        , ILogger<CronJobServiceBase> log,
        IServiceProvider serviceProvider) : base(hostedServiceSettings, log)
    {
        _serviceProvider = serviceProvider;
    }

    protected override async Task ExecuteTaskAsync(CancellationToken cancellationToken)
    {
        AppInsights.TrackTrace("Starting EventHubHostedService");
        _scope = _serviceProvider.CreateAsyncScope();
        _taskService = _scope.ServiceProvider.GetRequiredService<IEventServiceBusServiceTask>();
        await _taskService.StartAsync(cancellationToken);
    }

    protected override async Task DisposeScope()
    {
        await _taskService.StopAsync(CancellationToken.None);
        await _scope.DisposeAsync();
    }
}

public class EventServiceBusServiceTask : IEventServiceBusServiceTask
{
    private readonly IServiceBusManager _serviceBusManager;

    public EventServiceBusServiceTask(IServiceBusManager serviceBusManager)
    {
        _serviceBusManager = serviceBusManager;
    }

    public async Task StartAsync(CancellationToken cancellationToken)
    {
        await _serviceBusManager.StartProcessingAsync(cancellationToken);
    }

    public async Task StopAsync(CancellationToken cancellationToken)
    {
        await _serviceBusManager.StopProcessingAsync();
    }
}

问题根源分析

问题核心出在AddConsumerServiceBusConnection扩展方法中:你调用services.BuildServiceProvider()创建了临时服务提供实例,从这个实例中解析的ServiceBusManager是独立的临时对象,而非最终应用启动时使用的Singleton实例。

当应用最终构建主服务提供器时,会重新创建ServiceBusManager的Singleton实例,之前注册的处理器全部丢失,导致Hosted Service中获取的实例_processors为空。

解决方案

不要在服务配置阶段构建临时服务提供器,采用延迟注册或配置回调的方式确保消费者注册在主实例上执行:

方案一:利用Hosted Service延迟注册

修改扩展方法,通过临时Hosted Service在主服务提供器构建完成后执行消费者注册:

public static class ServiceCollectionExtensions
{
    public static IServiceCollection AddConsumerServiceBusConnection(
        this IServiceCollection services,
        Action<ServiceBusConnectionBuilder> configure)
    {
        services.AddHostedService(sp => 
        {
            var serviceBusManager = sp.GetRequiredService<IServiceBusManager>();
            var builder = new ServiceBusConnectionBuilder(serviceBusManager);
            configure(builder);
            return new ServiceBusConfigurationHostedService();
        });

        return services;
    }

    private class ServiceBusConfigurationHostedService : IHostedService
    {
        public Task StartAsync(CancellationToken cancellationToken) => Task.CompletedTask;
        public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
    }
}

方案二:通过构造函数注入配置动作

  1. 修改ServiceBusManager,添加配置动作集合并在初始化时执行:
public class ServiceBusManager : IServiceBusManager
{
    private readonly ServiceBusClient _serviceBusClient;
    private readonly IServiceProvider _serviceProvider;
    private readonly ConcurrentDictionary<string, ServiceBusProcessor> _processors;

    public ServiceBusManager(ServiceBusClient serviceBusClient,
        IServiceProvider serviceProvider,
        IOptions<ServiceBusSettings> serviceBusSettings,
        IEnumerable<Action<IServiceBusManager>> configurationActions)
    {
        _serviceBusClient = serviceBusClient ?? throw new ArgumentNullException(nameof(serviceBusClient));
        _serviceProvider = serviceProvider ?? throw new ArgumentNullException(nameof(serviceProvider));
        _processors = new ConcurrentDictionary<string, ServiceBusProcessor>();

        // 执行所有注册的消费者配置
        foreach (var action in configurationActions)
        {
            action(this);
        }
    }

    // 其余代码保持不变
}
  1. 修改扩展方法,注册配置动作:
public static class ServiceCollectionExtensions
{
    public static IServiceCollection AddConsumerServiceBusConnection(
        this IServiceCollection services,
        Action<ServiceBusConnectionBuilder> configure)
    {
        services.AddSingleton<Action<IServiceBusManager>>(sp => 
        {
            return manager => 
            {
                var builder = new ServiceBusConnectionBuilder(manager);
                configure(builder);
            };
        });

        return services;
    }
}

两种方案都能确保消费者注册在主ServiceBusManager实例上执行,Hosted Service启动时可获取到包含所有处理器的实例。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 22:05:54