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