【问题标题】:Multiple batch consumers throwing MessageLockLostException when receiving messages多个批处理消费者在接收消息时抛出 MessageLockLostException
【发布时间】:2021-11-12 21:51:22
【问题描述】:

我遇到了一个问题,我将消息发布到 Azure 服务总线主题。我有几个批量消费者订阅了这些主题(没有转发到队列)。

问题是有关于MessageLockLostExceptions的随机警告,感觉好像有问题,但不知道是什么问题。

我已将锁定持续时间设置为 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)

这是该问题的最小再现。它将设置所有内容并发布 50k 条消息。

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 服务总线实例的限制有关。第一张图显示了错误的发生,第二张图显示了 ASB 中受限制的请求数量。它们似乎相关性很好。这也可以解释为什么我无法可靠地重现错误。

【问题讨论】:

  • 如果这是一个简单的git clonedotnet run 示例,找时间看会容易得多。
  • @ChrisPatterson 你是对的!我把它放在这里(带有README 的说明):github.com/joelnotified/masstransit-batch-test。但是......我不能让它失败了。我认为我没有改变任何东西,所以我不知道发生了什么。将尝试看看我是否可以以更稳定的方式重现它,我会告诉你。
  • @ChrisPatterson 我一直在试图弄清楚为什么我不能可靠地重现它,并且看起来它与 ASB 的节流发生的时间有关。您认为:在节流期间是否会出现此类错误? ??????
  • 啊,在你的命名空间中节流肯定会导致一些问题。因为 ASB 本质上是在切断你的联系。奇怪的是,它们会使锁过期,除非您超过了锁计数阈值。可能是他们不做广告的那些“软”限制之一。

标签: azureservicebus masstransit azure-servicebus-topics azure-servicebus-subscriptions


【解决方案1】:

首先,我们需要确保我们是最新版本,如果不是,请升级它。

确保在 function.json 中添加重试策略,如下所示:

{
    "disabled": false,
    "bindings": [
        {
            ....
        }
    ],
    "retry": {
        "strategy": "fixedDelay",
        "maxRetryCount": 4,
        "delayInterval": "00:00:10"
    }
}

同时检查死信并在MS Docs 的帮助下进行配置。

我们需要在host.json中通过在JSON结构中包含参数“maxAutoLockRenewalDuration”:“00:05:00”来提及更新的锁定持续时间

下面是host.json的例子

{
    "version": "2.0",
    "extensions": {
        "serviceBus": {
            "clientRetryOptions":{
                "mode": "exponential",
                "tryTimeout": "00:01:00",
                "delay": "00:00:00.80",
                "maxDelay": "00:01:00",
                "maxRetries": 3
            },
            "prefetchCount": 0,
            "autoCompleteMessages": true,
            "maxAutoLockRenewalDuration": "00:05:00",
            "maxConcurrentCalls": 16,
            "maxConcurrentSessions": 8,
            "maxMessages": 1000,
            "sessionIdleTimeout": "00:01:00"
        }
    }
}

从上面的 JSON 中,将 autoCompleteMessages 的参数从 true 更改为 false 为 "autoCompleteMessages": false。

当设置为 false 时,您负责调用 MessageReceiver 方法来完成、放弃或死信消息。如果抛出异常(并且没有调用任何 MessageReceiver 方法),则锁定仍然存在。一旦锁过期,消息会重新排队,DeliveryCount 递增,锁会自动更新。

使用 ReceiveandLock 解决了这个问题issue

参考这些 SO 线程:SO1SO2(感谢回答的作者的详细解释)

【讨论】:

  • 这是一个 Azure Functions 答案,与问题没有密切关系。
猜你喜欢
  • 2020-11-11
  • 2013-08-25
  • 2015-01-02
  • 1970-01-01
  • 2013-08-29
  • 2016-06-08
  • 2016-05-04
  • 2023-03-24
  • 2022-08-13
相关资源
最近更新 更多