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

多批次消费者消费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异常,有哪些可行的优化解决思路?
异常发生时间统计
ASB限流请求统计


优化解决思路

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 20:24:03