多批次消费者消费Azure Service Bus消息触发MessageLockLostException异常
Azure Service Bus 限流触发 MessageLockLostException 异常优化问题
问题背景
我在将消息发布到Azure Service Bus主题时遇到问题:配置了多个批次消费者订阅这些主题(未配置队列转发),运行过程中随机抛出MessageLockLostException告警,暂时无法定位问题原因。
我已经将锁时长设置为5分钟,且报错几乎是立即触发的,因此可以排除锁超时的可能性。
报错示例
warn: MassTransit[0] Message Lock Lost: 5d5400005de20015b8d008d9a521105f Microsoft.Azure.ServiceBus.MessageLockLostException: The lock supplied is invalid. Either the lock expired, or the message has already been removed from the queue, or was received by a different receiver instance. at Microsoft.Azure.ServiceBus.Core.MessageReceiver.DisposeMessagesAsync(IEnumerable`1 lockTokens, Outcome outcome) at Microsoft.Azure.ServiceBus.RetryPolicy.RunOperation(Func`1 operation, TimeSpan operationTimeout) at Microsoft.Azure.ServiceBus.RetryPolicy.RunOperation(Func`1 operation, TimeSpan operationTimeout) at Microsoft.Azure.ServiceBus.Core.MessageReceiver.CompleteAsync(IEnumerable`1 lockTokens) at MassTransit.Transports.ReceivePipeDispatcher.Dispatch(ReceiveContext context, ReceiveLockContext receiveLock) at MassTransit.Transports.ReceivePipeDispatcher.Dispatch(ReceiveContext context, ReceiveLockContext receiveLock) at MassTransit.Transports.ReceivePipeDispatcher.Dispatch(ReceiveContext context, ReceiveLockContext receiveLock) at MassTransit.Azure.ServiceBus.Core.Transport.BrokeredMessageReceiver.MassTransit.Azure.ServiceBus.Core.Transport.IBrokeredMessageReceiver.Handle(Message message, CancellationToken cancellationToken, Action`1 contextCallback) warn: MassTransit[0] Message Lock Lost: 5d5400005de20015a5cd08d9a521105f Microsoft.Azure.ServiceBus.MessageLockLostException: The lock supplied is invalid. Either the lock expired, or the message has already been removed from the queue, or was received by a different receiver instance. at Microsoft.Azure.ServiceBus.Core.MessageReceiver.DisposeMessagesAsync(IEnumerable`1 lockTokens, Outcome outcome) at Microsoft.Azure.ServiceBus.RetryPolicy.RunOperation(Func`1 operation, TimeSpan operationTimeout) at Microsoft.Azure.ServiceBus.RetryPolicy.RunOperation(Func`1 operation, TimeSpan operationTimeout) at Microsoft.Azure.ServiceBus.Core.MessageReceiver.CompleteAsync(IEnumerable`1 lockTokens) at MassTransit.Transports.ReceivePipeDispatcher.Dispatch(ReceiveContext context, ReceiveLockContext receiveLock) at MassTransit.Transports.ReceivePipeDispatcher.Dispatch(ReceiveContext context, ReceiveLockContext receiveLock) at MassTransit.Transports.ReceivePipeDispatcher.Dispatch(ReceiveContext context, ReceiveLockContext receiveLock) at MassTransit.Azure.ServiceBus.Core.Transport.BrokeredMessageReceiver.MassTransit.Azure.ServiceBus.Core.Transport.IBrokeredMessageReceiver.Handle(Message message, CancellationToken cancellationToken, Action`1 contextCallback)
最小复现代码
csproj配置
<Project Sdk="Microsoft.NET.Sdk.Worker"> <PropertyGroup> <TargetFramework>net6.0</TargetFramework> <Nullable>enable</Nullable> <ImplicitUsings>enable</ImplicitUsings> <UserSecretsId>dotnet-WorkerService-C6197FFA-DCA6-4867-8576-A51ADAE04FD3</UserSecretsId> </PropertyGroup> <ItemGroup> <PackageReference Include="MassTransit" Version="7.2.3" /> <PackageReference Include="MassTransit.AspNetCore" Version="7.2.3" /> <PackageReference Include="MassTransit.Azure.ServiceBus.Core" Version="7.2.3" /> <PackageReference Include="MassTransit.EntityFrameworkCore" Version="7.2.3" /> <PackageReference Include="MassTransit.Prometheus" Version="7.2.3" /> <PackageReference Include="MassTransit.RabbitMQ" Version="7.2.3" /> <PackageReference Include="Microsoft.Extensions.Hosting" Version="6.0.0" /> </ItemGroup> </Project>
业务代码
using System; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using GreenPipes; using MassTransit; using MassTransit.Azure.ServiceBus.Core; using MassTransit.Topology; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using WorkerService; using IHost = Microsoft.Extensions.Hosting.IHost; IHost host = Host.CreateDefaultBuilder(args) .ConfigureServices(services => { const string connectionString = "<ASB ConnectionString here>"; Configure(services, connectionString); services.AddHostedService<Worker>(); }) .Build(); await host.RunAsync(); void Configure(IServiceCollection services, string connectionString) { services.AddMassTransit(busConfigurator => { busConfigurator.AddConsumer<TestConsumer1>(); busConfigurator.AddConsumer<TestConsumer2>(); busConfigurator.AddConsumer<TestConsumer3>(); busConfigurator.AddConsumer<TestConsumer4>(); busConfigurator.AddConsumer<TestConsumer5>(); busConfigurator.UsingAzureServiceBus((context, serviceBusBusFactoryConfigurator) => { serviceBusBusFactoryConfigurator.Host(connectionString); ConfigureSubsriptionEndpoint<TestConsumer1>(serviceBusBusFactoryConfigurator, context, "subscriber-1"); ConfigureSubsriptionEndpoint<TestConsumer2>(serviceBusBusFactoryConfigurator, context, "subscriber-2"); ConfigureSubsriptionEndpoint<TestConsumer3>(serviceBusBusFactoryConfigurator, context, "subscriber-3"); ConfigureSubsriptionEndpoint<TestConsumer4>(serviceBusBusFactoryConfigurator, context, "subscriber-4"); ConfigureSubsriptionEndpoint<TestConsumer5>(serviceBusBusFactoryConfigurator, context, "subscriber-5"); }); }); services.AddMassTransitHostedService(true); } void ConfigureSubsriptionEndpoint<TConsumer>(IServiceBusBusFactoryConfigurator serviceBusBusFactoryConfigurator, IBusRegistrationContext context, string subscriptionName) where TConsumer : class, IConsumer<Batch<IMyEvent>> { serviceBusBusFactoryConfigurator.SubscriptionEndpoint<IMyEvent>( subscriptionName, receiveEndpointConfigurator => { receiveEndpointConfigurator.LockDuration = TimeSpan.FromMinutes(5); receiveEndpointConfigurator.PublishFaults = false; receiveEndpointConfigurator.MaxAutoRenewDuration = TimeSpan.FromMinutes(30); receiveEndpointConfigurator.UseMessageRetry(r => r.Intervals(500, 2000)); receiveEndpointConfigurator.PrefetchCount = 1100; receiveEndpointConfigurator.ConfigureConsumer<TConsumer>( context, consumerConfigurator => { consumerConfigurator.Options<BatchOptions>(batchOptions => { batchOptions.MessageLimit = 100; batchOptions.TimeLimit = TimeSpan.FromSeconds(5); batchOptions.ConcurrencyLimit = 10; }); }); }); } namespace WorkerService { public class TestConsumer1 : IConsumer<Batch<IMyEvent>> { private readonly Random _random; private readonly ILogger<TestConsumer1> _logger; public TestConsumer1(ILogger<TestConsumer1> logger) { _logger = logger; _random = new Random(); } public async Task Consume(ConsumeContext<Batch<IMyEvent>> context) { _logger.LogInformation("{name} - Consuming {count}", nameof(TestConsumer1), context.Message.Length); await Task.Delay(TimeSpan.FromSeconds(_random.Next(4, 8))); } } public class TestConsumer2 : IConsumer<Batch<IMyEvent>> { private readonly Random _random; private readonly ILogger<TestConsumer2> _logger; public TestConsumer2(ILogger<TestConsumer2> logger) { _logger = logger; _random = new Random(); } public async Task Consume(ConsumeContext<Batch<IMyEvent>> context) { _logger.LogInformation("{name} - Consuming {count}", nameof(TestConsumer2), context.Message.Length); await Task.Delay(TimeSpan.FromSeconds(_random.Next(4, 8))); } } public class TestConsumer3 : IConsumer<Batch<IMyEvent>> { private readonly Random _random; private readonly ILogger<TestConsumer3> _logger; public TestConsumer3(ILogger<TestConsumer3> logger) { _logger = logger; _random = new Random(); } public async Task Consume(ConsumeContext<Batch<IMyEvent>> context) { _logger.LogInformation("{name} - Consuming {count}", nameof(TestConsumer3), context.Message.Length); await Task.Delay(TimeSpan.FromSeconds(_random.Next(4, 8))); } } public class TestConsumer4 : IConsumer<Batch<IMyEvent>> { private readonly Random _random; private readonly ILogger<TestConsumer4> _logger; public TestConsumer4(ILogger<TestConsumer4> logger) { _logger = logger; _random = new Random(); } public async Task Consume(ConsumeContext<Batch<IMyEvent>> context) { _logger.LogInformation("{name} - Consuming {count}", nameof(TestConsumer4), context.Message.Length); await Task.Delay(TimeSpan.FromSeconds(_random.Next(4, 8))); } } public class TestConsumer5 : IConsumer<Batch<IMyEvent>> { private readonly Random _random; private readonly ILogger<TestConsumer5> _logger; public TestConsumer5(ILogger<TestConsumer5> logger) { _logger = logger; _random = new Random(); } public async Task Consume(ConsumeContext<Batch<IMyEvent>> context) { _logger.LogInformation("{name} - Consuming {count}", nameof(TestConsumer5), context.Message.Length); await Task.Delay(TimeSpan.FromSeconds(_random.Next(4, 8))); } } [EntityName("my-event")] public interface IMyEvent { } public class Worker : BackgroundService { private readonly ILogger<Worker> _logger; private readonly IBus _bus; public Worker( ILogger<Worker> logger, IBus bus) { _logger = logger; _bus = bus; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("Worker running at: {time}", DateTimeOffset.Now); var tasks = new List<Task>(); var count = 50000; for (int i = 0; i < count; i++) { tasks.Add(_bus.Publish<IMyEvent>(new { })); } await Task.WhenAll(tasks); } } }
问题更新
已确认该报错与Azure Service Bus实例的限流行为高度相关,报错发生的时间点与ASB的限流请求统计数据完全吻合,这也解释了此前无法稳定复现问题的原因。请问针对该类由限流触发的MessageLockLostException异常,有哪些可行的优化解决思路?

优化解决思路
1. ASB资源层优化
- 升级服务层级:如果使用的是标准层ASB,升级到高级层,可获得更高的吞吐量阈值、更低的延迟和更稳定的限流策略。
- 调整吞吐量单位:高级层ASB可动态调整吞吐量单位(TU),根据峰值流量按需扩容,避免固定配额被打满触发限流。
- 开启主题分区:启用ASB主题的分区功能,消息会分散存储在多个分片上,分摊流量压力,降低单分片限流概率。
- 就近部署:确保服务和ASB实例在同一Azure区域,减少跨区域请求延迟导致的额外超时和重试请求。
2. 消费端配置优化
- 降低预取数量:当前
PrefetchCount = 1100配置过高,限流场景下预取的大量消息来不及在锁有效期内完成确认,会直接触发锁丢失。建议调整为批量消息限制的23倍,例如当前批次上限是100,可调整为200300。 - 调整批量消费参数:降低
ConcurrencyLimit并发批量数,同时可适当下调MessageLimit单批次消息上限,减少单批次处理时间和确认请求的压力。 - 优化重试策略:将普通间隔重试改为针对限流错误的指数退避重试,单独捕获
ServiceBusException并判断IsTransient和状态码是否为429,针对性配置更长的重试间隔,避免无效重试加重限流。 - 开启批量确认:配置MassTransit批量确认消息的能力,减少向ASB发送Complete请求的次数,降低请求密度。
3. 发布端优化
- 控制发布并发:当前代码一次性并发发布5万条消息,突发流量直接打满ASB限流阈值。可通过
SemaphoreSlim限制并发发布数量,例如限制同时最多100个发布请求。 - 批量发布:使用ASB的批量发布接口,单次请求最多发送100条消息,大幅减少请求数,降低限流概率。
- 削峰填谷:如果业务允许,可将突发流量错峰发布,避免短时间内请求量突增。
4. 兜底容错配置
- 配置死信队列:将锁丢失的异常消息自动转入死信队列,避免消息丢失,后续可单独消费死信队列做补处理。
- 保证消费幂等:锁丢失后ASB会自动重投消息,需要确保消费逻辑的幂等性,避免重复消费导致业务异常。
内容的提问来源于stack exchange,提问作者Joel
相关产品推荐
相关产品推荐

