.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; } }
方案二:通过构造函数注入配置动作
- 修改
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); } } // 其余代码保持不变 }
- 修改扩展方法,注册配置动作:
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
相关产品推荐
相关产品推荐

