MassTransit-Azure Service Bus消费者重复消费同一消息问题求助
问题
我有一个基于MassTransit-Azure Service Bus的Worker服务与API服务,API作为消息发送端,Worker作为接收端。在Worker服务中创建了消费者处理队列消息,但处理过程中MassTransit会自动重新初始化消费者并再次触发Consume方法处理同一消息。
需要实现:同一消息仅由一个消费者实例获取并处理,其他同类型消费者实例(部署在不同服务器)在该消息处理期间无法获取它。另外,不清楚如何在MassTransit中启用Azure的ReceiveAndDelete模式(区别于Peek-Mode)。
相关配置代码如下:
API配置
public static IServiceCollection AddMessageBusConfig(this IServiceCollection services, IDictionary<string, string> messageBusOptions) { var queuePrefix = messageBusOptions["QueuePrefix"]; return services.AddMassTransit(x => { EndpointConvention.Map<MigrateAccountancy.Command>(new Uri($"queue:{queuePrefix}-migrate-accountancies-event")); x.UsingAzureServiceBus((context, cfg) => { cfg.Host(messageBusOptions["ConnectionString"]); cfg.ConfigureEndpoints(context); }); }); }
Worker服务配置
public static IServiceCollection AddMessageBus(this IServiceCollection services, IDictionary<string, string> messageBusOptions) { var queuePrefix = messageBusOptions["QueuePrefix"]; return services.AddMassTransit(x => { x.AddConsumer<MigrateAccountancyConsumer>(); x.UsingAzureServiceBus((context, cfg) => { cfg.Host(messageBusOptions["ConnectionString"]); cfg.ReceiveEndpoint($"{queuePrefix}-migrate-accountancies-event", e => { // skip move message to _error queue on exception since errors are logged in Application Insight e.DiscardFaultedMessages(); e.ConfigureConsumer<MigrateAccountancyConsumer>(context, c => { }); }); }); }); }
消费者代码
public class MigrateAccountancyConsumer : IConsumer<MigrateAccountancy.Command> { private readonly IMediator _mediator; public MigrateAccountancyConsumer(IMediator mediator) { _mediator = mediator; } public async Task Consume(ConsumeContext<MigrateAccountancy.Command> context) { await _mediator.Send(context.Message); } }
解决方案
一、阻止同一消息被重复消费的配置
MassTransit默认借助Azure Service Bus的锁机制保证同一消息仅被一个消费者实例处理,出现重复触发Consume的情况,大多是因为消息锁超时,或者消费者处理时长超过锁有效期,导致Service Bus判定消息未处理完成,重新将其放回队列。
你可以通过调整锁相关参数解决该问题,在Worker的ReceiveEndpoint配置中添加以下设置:
cfg.ReceiveEndpoint($"{queuePrefix}-migrate-accountancies-event", e => { e.DiscardFaultedMessages(); // 根据业务处理时长设置消息锁有效期,单位为秒/分钟 e.LockDuration = TimeSpan.FromMinutes(5); // 开启自动续期,续期间隔需小于LockDuration,避免锁过期 e.AutoRenewTimeout = TimeSpan.FromMinutes(4); // 设置最大交付次数,防止消息无限重复 e.MaxDeliveryCount = 3; e.ConfigureConsumer<MigrateAccountancyConsumer>(context); });
同时要确保消费者处理逻辑没有未捕获的异常——异常会触发MassTransit将消息重新放回队列,即使设置了DiscardFaultedMessages也需要确认其生效逻辑是否符合预期。
二、启用Azure Service Bus的ReceiveAndDelete模式
MassTransit默认使用Peek-Lock模式(先锁定消息,处理完成后再删除),若要切换到ReceiveAndDelete模式,可在ReceiveEndpoint配置中设置ReceiveMode:
cfg.ReceiveEndpoint($"{queuePrefix}-migrate-accountancies-event", e => { e.DiscardFaultedMessages(); // 启用ReceiveAndDelete模式 e.ReceiveMode = Azure.Messaging.ServiceBus.ReceiveMode.ReceiveAndDelete; e.ConfigureConsumer<MigrateAccountancyConsumer>(context); });
注意:ReceiveAndDelete模式下,消息被消费者接收后立即删除,若处理失败消息会直接丢失,仅适合允许消息丢失或有额外容错机制的场景。
修改后的Worker完整配置代码:
public static IServiceCollection AddMessageBus(this IServiceCollection services, IDictionary<string, string> messageBusOptions) { var queuePrefix = messageBusOptions["QueuePrefix"]; return services.AddMassTransit(x => { x.AddConsumer<MigrateAccountancyConsumer>(); x.UsingAzureServiceBus((context, cfg) => { cfg.Host(messageBusOptions["ConnectionString"]); cfg.ReceiveEndpoint($"{queuePrefix}-migrate-accountancies-event", e => { // skip move message to _error queue on exception since errors are logged in Application Insight e.DiscardFaultedMessages(); // 配置锁参数防止重复消费 e.LockDuration = TimeSpan.FromMinutes(5); e.AutoRenewTimeout = TimeSpan.FromMinutes(4); e.MaxDeliveryCount = 3; // 启用ReceiveAndDelete模式(按需选择) // e.ReceiveMode = Azure.Messaging.ServiceBus.ReceiveMode.ReceiveAndDelete; e.ConfigureConsumer<MigrateAccountancyConsumer>(context); }); }); }); }
内容的提问来源于stack exchange,提问作者Kiet Duong

