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

.NET Web应用托管服务接收Azure Service Bus消息不全问题求助

问题:Azure Service Bus主题消息无法全部接收,批量发送部分消息进入死信队列

我在.NET Web应用中搭建了托管服务,用于接收Azure Service Bus主题的消息。目前遇到的问题是无法接收全部消息:批量并发发送时,仅能接收随机数量的消息(例如20条中仅12条被接收),其余消息进入死信队列。

已尝试的解决步骤

  • 调大最大并发调用数,有一定效果但无法保证全部接收
  • 添加预取计数

测试情况

通过Azure门户的Service Bus资源功能测试发送消息:

  • 发送500条无间隔消息时无法全部接收
  • 设置1秒间隔发送500条时,所有消息均被接收

我需要构建可靠的事件驱动架构,不能让消息处理的完整性成为随机事件,希望解决该问题。


相关代码

Startup.cs

...
public void ConfigureServices(IServiceCollection services)
{
    services.AddSingleton<IServiceBusTopicSubscription, ServiceBusSubscription>();
    services.AddHostedService<WorkerServiceBus>();
}
...

WorkerService.cs

public class WorkerServiceBus : IHostedService, IDisposable
{
    private readonly ILogger<WorkerServiceBus> _logger;
    private readonly IServiceBusTopicSubscription _serviceBusTopicSubscription;

    public WorkerServiceBus(IServiceBusTopicSubscription serviceBusTopicSubscription,
        ILogger<WorkerServiceBus> logger)
    {
        _serviceBusTopicSubscription = serviceBusTopicSubscription;
        _logger = logger;
    }

    public async Task StartAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Starting the service bus queue consumer and the subscription");
        await _serviceBusTopicSubscription.PrepareFiltersAndHandleMessages().ConfigureAwait(false);
    }

    public async Task StopAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Stopping the service bus queue consumer and the subscription");
        await _serviceBusTopicSubscription.CloseSubscriptionAsync().ConfigureAwait(false);
    }

    public void Dispose()
    {
        Dispose(true);
        GC.SuppressFinalize(this);
    }

    protected virtual async void Dispose(bool disposing)
    {
        if (disposing)
        {
            await _serviceBusTopicSubscription.DisposeAsync().ConfigureAwait(false);
        }
    }
}

ServiceBusSubscription.cs

public class ServiceBusSubscription : IServiceBusTopicSubscription
{
    private readonly IConfiguration _configuration;
    private const string TOPIC_PATH = "test";
    private const string SUBSCRIPTION_NAME = "test-subscriber";
    private readonly ILogger _logger;
    private readonly ServiceBusClient _client;
    private readonly IServiceScopeFactory _scopeFactory;
    private ServiceBusProcessor _processor;

    public ServiceBusBookingsSubscription(IConfiguration configuration,
        ILogger<ServiceBusBookingsSubscription> logger,
        IServiceScopeFactory scopeFactory)
    {
        _configuration = configuration;
        _logger = logger;
        _scopeFactory = scopeFactory;

        var connectionString = _configuration.GetConnectionString("ServiceBus");
        var serviceBusOptions = new ServiceBusClientOptions()
        {
            TransportType = ServiceBusTransportType.AmqpWebSockets
        };
        _client = new ServiceBusClient(connectionString, serviceBusOptions);
    }

    public async Task PrepareFiltersAndHandleMessages()
    {
        ServiceBusProcessorOptions _serviceBusProcessorOptions = new ServiceBusProcessorOptions
        {
            MaxConcurrentCalls = 200,
            AutoCompleteMessages = false,
            PrefetchCount = 1000,
        };

        _processor = _client.CreateProcessor(TOPIC_PATH, SUBSCRIPTION_NAME, _serviceBusProcessorOptions);
        _processor.ProcessMessageAsync += ProcessMessagesAsync;
        _processor.ProcessErrorAsync += ProcessErrorAsync;

        await _processor.StartProcessingAsync().ConfigureAwait(false);
    }

    private async Task ProcessMessagesAsync(ProcessMessageEventArgs args)
    {
        _logger.LogInformation("Received message from service bus");
        _logger.LogInformation("Message: {args.Message.Body}");
        var payload = args.Message.Body.ToObjectFromJson<List<SchedulerBookingViewModel>>();

        // Create scoped dbcontext
        using var scope = _scopeFactory.CreateScope();
        var dbContext = scope.ServiceProvider.GetRequiredService<dbContext>();

        // Process payload
        await new TestServiceBus().DoThings(payload);

        await args.CompleteMessageAsync(args.Message).ConfigureAwait(false);
    }

    private Task ProcessErrorAsync(ProcessErrorEventArgs arg)
    {
        _logger.LogError(arg.Exception, "Message handler encountered an exception");
        _logger.LogError("- ErrorSource: {arg.ErrorSource}");
        _logger.LogError("- Entity Path: {arg.EntityPath}");
        _logger.LogError("- FullyQualifiedNamespace: {arg.FullyQualifiedNamespace}");

        return Task.CompletedTask;
    }

    public async ValueTask DisposeAsync()
    {
        if (_processor != null)
        {
            await _processor.DisposeAsync().ConfigureAwait(false);
        }

        if (_client != null)
        {
            await _client.DisposeAsync().ConfigureAwait(false);
        }
    }

    public async Task CloseSubscriptionAsync()
    {
        await _processor.CloseAsync().ConfigureAwait(false);
    }
}

解决方案

1. 修复异常处理逻辑

当前代码未捕获处理过程中的异常,一旦消息处理失败(如数据库连接超时、业务逻辑报错),会导致消息被自动重试多次后进入死信。需添加异常捕获,根据异常类型决定重试或死信:

private async Task ProcessMessagesAsync(ProcessMessageEventArgs args)
{
    var messageId = args.Message.MessageId;
    _logger.LogInformation("Received message {MessageId} from service bus", messageId);
    try
    {
        var payload = args.Message.Body.ToObjectFromJson<List<SchedulerBookingViewModel>>();

        using var scope = _scopeFactory.CreateScope();
        var dbContext = scope.ServiceProvider.GetRequiredService<dbContext>();

        await new TestServiceBus().DoThings(payload);

        await args.CompleteMessageAsync(args.Message).ConfigureAwait(false);
        _logger.LogInformation("Completed message {MessageId}", messageId);
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Failed to process message {MessageId}", messageId);
        // 可恢复异常(如数据库连接异常)延迟重试
        if (ex is DbUpdateException || ex is SqlException)
        {
            await Task.Delay(TimeSpan.FromSeconds(2));
            await args.AbandonMessageAsync(args.Message).ConfigureAwait(false);
        }
        else
        {
            // 不可恢复异常直接死信,并记录原因
            await args.DeadLetterMessageAsync(args.Message, "ProcessingFailed", ex.Message).ConfigureAwait(false);
        }
    }
}

2. 调整并发与资源配置

  • 降低MaxConcurrentCalls:当前设置为200,可能超出下游资源(如数据库连接池)承载能力,建议根据实际资源情况调整为50~100。
  • 配置全局重试策略:在ServiceBusClientOptions中添加重试规则,避免消息快速进入死信:
var serviceBusOptions = new ServiceBusClientOptions()
{
    TransportType = ServiceBusTransportType.AmqpWebSockets,
    RetryOptions = new ServiceBusRetryOptions
    {
        Mode = ServiceBusRetryMode.Exponential,
        MaxRetries = 3,
        Delay = TimeSpan.FromSeconds(1),
        MaxDelay = TimeSpan.FromSeconds(5)
    }
};
  • 扩容数据库连接池:在数据库连接字符串中增加Max Pool Size,避免并发处理时连接耗尽:
"ConnectionStrings": {
    "ServiceBus": "...",
    "YourDbConnection": "Server=...;Database=...;Trusted_Connection=True;Max Pool Size=100;"
}

3. 完善日志与监控

  • 增加消息ID、处理时长等关键信息的日志,方便定位具体失败的消息。
  • 通过Azure Monitor设置死信队列的告警规则,及时发现异常。
  • 检查订阅的自定义死信规则,确认是否有额外规则导致消息被移入死信。

4. 优化发送与处理逻辑

  • 避免无间隔批量发送:使用ServiceBus的批量发送API控制发送速率,或添加固定间隔。
  • 拆分大消息:如果单条消息Payload过大,拆分消息减少单次处理时长,避免超时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 04:30:57