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

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消费者接收。

  1. 定义业务异常类:
public class BusinessException : Exception
{
    public BusinessException(string message) : base(message) { }
}
  1. 修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 00:27:22