.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
相关产品推荐
相关产品推荐

