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

.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 08:05:55