Azure Service Bus+MassTransit消息异常重复调用问题咨询
问题原因分析
Azure Functions的ServiceBusTrigger默认行为是:当函数执行抛出异常时,会将消息放弃(Abandon),Azure Service Bus会自动将消息重新投递回订阅,直到达到默认的最大投递次数(10次)后移入死信队列。而Service1中使用MassTransit托管消费者时,MassTransit默认会在消费者抛出异常时直接标记消息失败并发布Fault消息,不会触发Service Bus的重试机制——两者的差异源于消息异常处理的责任主体不同:前者由Azure Functions Runtime管控,后者由MassTransit框架接管。
解决方案
方案1:区分异常类型,手动控制消息状态
自定义业务异常标记类,在Function中捕获业务异常时直接完成消息(避免重试),仅系统异常触发Service Bus重试,同时可手动发布Fault消息让Service1的Fault消费者接收。
- 定义业务异常类:
public class BusinessException : Exception { public BusinessException(string message) : base(message) { } }
- 修改Function执行代码:
public async Task Run([ServiceBusTrigger(TopicName, SubscriptionName, Connection = Connection)] ServiceBusReceivedMessage message, ServiceBusMessageActions messageActions, IBus bus, CancellationToken cancellationToken) { try { await _receiver.HandleConsumer<TestMessageConsumer>(TopicName, SubscriptionName, message, cancellationToken); await messageActions.CompleteMessage(message, cancellationToken); } catch (BusinessException ex) { // 业务异常:直接完成消息,不重试 await messageActions.CompleteMessage(message, cancellationToken); // 手动发布Fault消息,让Service1的Fault消费者触发 await bus.Publish<Fault<TestMessage>>(new { MessageId = message.MessageId, ExceptionInfo = new[] { new ExceptionInfo { Message = ex.Message, StackTrace = ex.StackTrace } }, Timestamp = DateTime.UtcNow }, cancellationToken); } catch (Exception) { // 系统异常:放弃消息,触发Service Bus重试 await messageActions.AbandonMessage(message, cancellationToken); } }
方案2:通过MassTransit配置接管异常处理
在Service2的MassTransit配置中,为消费者添加错误处理策略,让框架区分异常类型并控制重试逻辑,无需依赖Azure Functions Runtime的默认行为。
修改Service2的MassTransit配置:
services.AddMassTransitForAzureFunctions(cfg => { cfg.AddConsumer<TestMessageConsumer>(c => { // 业务异常直接发布Fault,跳过重试 c.AddFault(options => options.SkipRetry = true); // 仅对系统异常进行有限次数重试 c.UseMessageRetry(r => { r.Handle<TimeoutException, IOException>(); // 指定需要重试的系统异常类型 r.MaximumAttempts(3); // 最多重试3次 r.Interval(1, TimeSpan.FromSeconds(2)); // 重试间隔 }); }); }, "ServiceBusConnection", (context, configuration) => { var connectionStrings = context.GetRequiredService<IOptions<ConnectionStringOption>>(); var settings = new HostSettings { ServiceUri = new Uri("sb://" + connectionStrings.Value.Messaging), TokenCredential = new DefaultAzureCredential(), }; configuration.Host(settings); configuration.ConfigureEndpoints(context); });
同时简化Function代码,让MassTransit完全接管消息处理:
public async Task Run([ServiceBusTrigger(TopicName, SubscriptionName, Connection = Connection)] ServiceBusReceivedMessage message, CancellationToken cancellationToken) { await _receiver.HandleConsumer<TestMessageConsumer>(TopicName, SubscriptionName, message, cancellationToken); }
方案3:全局调整订阅的重试参数(不推荐单独使用)
通过Service Bus管理客户端修改订阅的最大投递次数,全局控制重试次数,但无法区分业务/系统异常,适合统一限制重试的场景:
修改Service1中的订阅创建代码:
if (!await adminClient.SubscriptionExistsAsync(formatter.MessageName<TestMessage>(), formatterCheckFunc.ConsumerFromMessageName<TestMessageConsumer>())) { var subscriptionOptions = new CreateSubscriptionOptions(formatter.MessageName<TestMessage>(), formatterCheckFunc.ConsumerFromMessageName<TestMessageConsumer>()) { MaxDeliveryCount = 1 // 仅投递1次,异常后直接进入死信 }; await adminClient.CreateSubscriptionAsync(subscriptionOptions); }
内容的提问来源于stack exchange,提问作者Petr Klekner
相关产品推荐
相关产品推荐

