.NET 6应用Azure Service Bus接收器因闲置超时断开求助
问题描述
我有一个基础的.NET 6应用,仅用于接收并处理Azure Service Bus消息,正常运行时符合预期业务逻辑。但在接收并处理完消息后,Azure Service Bus接收器会进入闲置状态,导致链接被关闭。
日志信息
应用为Docker镜像,通过Jenkins/ArgoCD部署在AKS上,日志如下:
info: Inventory_Create_Worker.API.Consumer.CreateInventoryConsumerService[0] Message received and handled: {"Title":"inventory creation event","Description":"Extract from excel file","Event":":CONFIDENTIAL_INVENTORY_DATA"} info: Azure.Messaging.ServiceBus[38] Receive Link Closed. Identifier: queue-inventory-creation-dev-xxxxxx-xxxx-xxxx-xxxx-xxxxxxxxx, SessionId: , linkException: Azure.Messaging.ServiceBus.ServiceBusException: The link 'amqps://personaldevservicebus-development.servicebus.windows.net/-f1a43db3;0:5:6:source(address:/queue-inventory-creation-dev,filter:[])' is force detached. Code: aggregate-link1318524517. Details: AmqpMessagePartitioningEntityMessageCache.IdleTimerExpired: Idle timeout in Seconds: 900. (GeneralError).
IOC注入代码
services.AddAzureClients(clientsBuilder => { clientsBuilder.AddServiceBusClient(connectionString:azureQueueConfig.ConnectionString) .WithName(ConfirmationCreateServiceBusClientName) .ConfigureOptions(options => { options.RetryOptions.Delay = TimeSpan.FromSeconds(2); options.RetryOptions.MaxDelay = TimeSpan.FromMinutes(2); options.ConnectionIdleTimeout = TimeSpan.FromMinutes(5); options.RetryOptions.MaxRetries = 5; }); });
最初未配置任何.ConfigureOptions,怀疑问题出在此配置中。
后台服务代码
using Shared; using Shared.AzureServiceBus; using Shared.Interfaces; using Shared.Interfaces.Infrastructure; using Shared.Interfaces.Provider; using Shared.Interfaces.Service; using Shared.Models.Config; using Shared.RabbitMq; using Shared.Services; using Inventory_Create_Worker.Application.Common.Interfaces.Authentication; using Inventory_Create_Worker.Infrastructure.Authentication; using Inventory_Create_Worker.Shared.Constants; using Microsoft.Extensions.Azure; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using static Inventory_Create_Worker.Shared.Constants.AzureConstant; public class CreateInventoryConsumerService : BackgroundService { private readonly ILogger<CreateInventoryConsumerService> _logger; private ServiceBusReceiver _serviceBusReceiver; private ServiceBusClient _serviceBusClient; private readonly IAzureClientFactory<ServiceBusClient> _azureClientFactory; private readonly AzureQueueConfig _azureQueueConfig; private readonly IServiceProvider _serviceProvider; public CreateInventoryConsumerService(IAzureClientFactory<ServiceBusClient> serviceBusClientFactory, ILogger<CreateInventoryConsumerService> logger, IOptions<AzureQueueConfig> azureQueueOptions, IServiceProvider serviceProvider) { _logger = logger; _serviceProvider = serviceProvider; _azureClientFactory = serviceBusClientFactory; _azureQueueConfig = azureQueueOptions.Value; _serviceBusClient = _azureClientFactory.CreateClient(name: AzureConstant.InventoryCreateServiceBusClientName); _serviceBusReceiver = _serviceBusClient.CreateReceiver(queueName: _azureQueueConfig.QueueName); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { try { var receivedMessage = await _serviceBusReceiver.ReceiveMessageAsync(cancellationToken: stoppingToken); if (receivedMessage != null) { using (var scope = _serviceProvider.CreateScope()) { var serviceProvider = scope.ServiceProvider; var mediator = serviceProvider.GetRequiredService<IMediator>(); Console.WriteLine("BEGINNING TO CONSUME SERVICE"); var createInventoryEvent = Encoding.UTF8.GetString(receivedMessage.Body); Console.WriteLine("var createInventoryEvent = Encoding.UTF8.GetString DONE"); CreateInventoryFromEventCommand command = new(createInventoryEvent); Console.WriteLine( "CreateInventoryFromEventCommand command = new(createInventoryEvent); DONE"); await mediator.Send(command, stoppingToken); Console.WriteLine("_mediator.Send(command) DONE"); _logger.LogInformation("Message received and handled: " + createInventoryEvent); Console.WriteLine("Message received and handled: " + createInventoryEvent); } } } catch (Exception ex) { _logger.LogError($"Message could not be executed in Inventory_Create_Worker.CreateInventoryConsumerService.ExecuteAsync, Exception raised: {ex.Message} )"); Console.WriteLine($"Message could not be executed in Inventory_Create_Worker.CreateInventoryConsumerService.ExecuteAsync, Exception raised: {ex.Message}"); RecreateServiceBusInstances(); } } /// <summary> /// When our service bus crashes, we need to recreate it. /// If not we are not listening to Azure Service bus events anymore and our Inventory consumer service is not guaranteed anymore /// </summary> private void RecreateServiceBusInstances() { try { _logger.LogInformation("Recreating a new ServiceBusReceiver and a new ServiceBusClient instance"); Console.WriteLine("Recreating a new ServiceBusReceiver and a new ServiceBusClient instance"); _serviceBusClient = _azureClientFactory.CreateClient(name: AzureConstant.InventoryCreateServiceBusClientName); _serviceBusReceiver = _serviceBusClient.CreateReceiver(queueName: _azureQueueConfig.QueueName); } catch (Exception ex) { _logger.LogError($"Failed to recreate ServiceBus instances: {ex.Message}"); Console.WriteLine($"Failed to recreate ServiceBus instances: {ex.Message}"); } } }
解决方案
核心问题分析
当前代码仅调用一次ReceiveMessageAsync,处理完消息后后台服务直接进入闲置状态,没有持续监听队列,导致Service Bus连接因长时间无活动触发闲置超时(日志显示超时时间为900秒)。另外,配置中的ConnectionIdleTimeout设为5分钟,比Service Bus默认闲置超时短,会加速连接关闭。
修复步骤
1. 改为持续监听队列
在ExecuteAsync中添加循环,确保处理完一条消息后继续监听下一条,避免服务闲置:
protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { var receivedMessage = await _serviceBusReceiver.ReceiveMessageAsync( maxWaitTime: TimeSpan.FromMinutes(1), // 设置合理等待时间,避免无限阻塞 cancellationToken: stoppingToken); if (receivedMessage != null) { using (var scope = _serviceProvider.CreateScope()) { var serviceProvider = scope.ServiceProvider; var mediator = serviceProvider.GetRequiredService<IMediator>(); var createInventoryEvent = Encoding.UTF8.GetString(receivedMessage.Body); CreateInventoryFromEventCommand command = new(createInventoryEvent); await mediator.Send(command, stoppingToken); _logger.LogInformation("Message received and handled: " + createInventoryEvent); // 处理完成后手动完成消息,避免重复消费 await _serviceBusReceiver.CompleteMessageAsync(receivedMessage, stoppingToken); } } } catch (ServiceBusException ex) when (ex.Reason == ServiceBusFailureReason.LinkDetached) { _logger.LogWarning("Service Bus链接断开,正在重建接收器..."); RecreateServiceBusInstances(); } catch (Exception ex) { _logger.LogError($"处理消息出错: {ex.Message}"); // 可根据需求添加死信队列逻辑 } } }
2. 调整连接闲置超时配置
将ConnectionIdleTimeout设置为大于等于Service Bus默认闲置超时(900秒=15分钟),或者直接移除该配置使用默认值:
services.AddAzureClients(clientsBuilder => { clientsBuilder.AddServiceBusClient(connectionString:azureQueueConfig.ConnectionString) .WithName(ConfirmationCreateServiceBusClientName) .ConfigureOptions(options => { options.RetryOptions.Delay = TimeSpan.FromSeconds(2); options.RetryOptions.MaxDelay = TimeSpan.FromMinutes(2); options.RetryOptions.MaxRetries = 5; // 移除或调整为15分钟以上 // options.ConnectionIdleTimeout = TimeSpan.FromMinutes(15); }); });
3. 推荐使用ServiceBusProcessor替代手动Receiver
ServiceBusProcessor是官方推荐的消息消费方式,内置自动连接恢复、批量处理等功能,比手动管理ServiceBusReceiver更稳定:
注册Processor(IOC容器中)
services.AddSingleton<ServiceBusProcessor>(sp => { var clientFactory = sp.GetRequiredService<IAzureClientFactory<ServiceBusClient>>(); var config = sp.GetRequiredService<IOptions<AzureQueueConfig>>().Value; var client = clientFactory.CreateClient(AzureConstant.InventoryCreateServiceBusClientName); return client.CreateProcessor(config.QueueName, new ServiceBusProcessorOptions { MaxAutoLockRenewalDuration = TimeSpan.FromMinutes(5), // 防止长时间处理时消息锁过期 AutoCompleteMessages = false // 手动控制消息完成 }); });
重构后台服务使用Processor
public class CreateInventoryConsumerService : BackgroundService { private readonly ILogger<CreateInventoryConsumerService> _logger; private readonly ServiceBusProcessor _processor; private readonly IServiceProvider _serviceProvider; public CreateInventoryConsumerService(ServiceBusProcessor processor, ILogger<CreateInventoryConsumerService> logger, IServiceProvider serviceProvider) { _processor = processor; _logger = logger; _serviceProvider = serviceProvider; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { async Task ProcessMessageAsync(ProcessMessageEventArgs args) { using var scope = _serviceProvider.CreateScope(); var mediator = scope.ServiceProvider.GetRequiredService<IMediator>(); var createInventoryEvent = Encoding.UTF8.GetString(args.Message.Body); var command = new CreateInventoryFromEventCommand(createInventoryEvent); await mediator.Send(command, stoppingToken); _logger.LogInformation("Message received and handled: " + createInventoryEvent); await args.CompleteMessageAsync(args.Message, stoppingToken); } Task ProcessErrorAsync(ProcessErrorEventArgs args) { _logger.LogError($"处理消息出错: {args.Exception.Message}"); return Task.CompletedTask; } _processor.ProcessMessageAsync += ProcessMessageAsync; _processor.ProcessErrorAsync += ProcessErrorAsync; await _processor.StartProcessingAsync(stoppingToken); await stoppingToken.WhenCanceled(); await _processor.StopProcessingAsync(stoppingToken); } }
额外注意事项
- 必须确保消息处理完成后调用
CompleteMessageAsync,避免消息重新入队重复消费 - 对于长时间运行的消息,设置合理的
MaxAutoLockRenewalDuration,防止消息锁过期 - 使用依赖注入管理
ServiceBusClient和ServiceBusProcessor,确保单例复用,避免频繁创建连接
内容的提问来源于stack exchange,提问作者Dodge
相关产品推荐
相关产品推荐

