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

使用MassTransit+Azure Service Bus时,消费者抛异常后停止消费

MassTransit搭配Azure Service Bus时,消费者抛出异常后停止消费新消息

问题描述

我用MassTransit结合Azure Service Bus,正常运行没问题,但消费者抛出异常时会出问题:消息会在Azure门户的订阅概览中变成死信,之后消费者就停止消费新消息了。

配置代码

_azureServiceBus = Bus.Factory.CreateUsingAzureServiceBus(cfg =>
{
    cfg.Host(new HostSettings
    {
        ServiceUri = new Uri(pharmacyEndpointSettings.Endpoint.BdcpServiceBus),
        TransportType = ServiceBusTransportType.AmqpTcp,
        TokenCredential = new ClientSecretCredential(
            pharmacyEndpointSettings.Endpoint.BdcpServiceTenantId,
            pharmacyEndpointSettings.Endpoint.BdcpServiceClientId,
            pharmacyEndpointSettings.Endpoint.BdcpServiceClientSecret)
    });
    
    cfg.UseMessageRetry(c => c.Immediate(5));
    cfg.AutoDeleteOnIdle = TimeSpan.FromDays(AutoDeleteSubscriptionOnIdleInDays);
    cfg.DefaultMessageTimeToLive = TimeSpan.FromDays(MessageTimeToLiveInDays);
    cfg.EnableDeadLetteringOnMessageExpiration = true;

    cfg.Message<MyMessage>(m => m.SetEntityName("MyTopicName"));
    cfg.SubscriptionEndpoint<T>(_subscriptionName.Value, se =>
    {
        se.Rule = new CreateRuleOptions("Receiver", new SqlRuleFilter("receiver='all' OR receiver='MySpecificReceiver'"));
        
        se.Consumer<MyMessage>();
        se.ConfigureDeadLetterQueueDeadLetterTransport();
        se.ConfigureDeadLetterQueueErrorTransport();
    });
}

_azureServiceBus.Start();

消费者代码

public class MyMessageConsumer : IConsumer<T>
{
   public override Task Consume(ConsumeContext<MyMessage> context)
   {
        if (context.Message.Id == "INVALID")
        {
            return Task.FromException(new Exception("invalid id"));
        }

        return Task.CompletedTask;
   }
}

死信中的OperationCancelledException堆栈

à GreenPipes.Internals.Extensions.TaskExtensions.<>c__DisplayClass5_01.<<OrCanceled>g__WaitAsync|0>d.MoveNext() --- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée --- à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task) à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task) à System.Runtime.CompilerServices.ConfiguredTaskAwaitable1.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Azure.ServiceBus.Core.Pipeline.SendEndpointContextFactory.d__7.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable1.ConfiguredTaskAwaiter.GetResult() à GreenPipes.Agents.PipeContextSupervisor1.d__7.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw()
à GreenPipes.Agents.PipeContextSupervisor1.<GreenPipes-IPipeContextSource<TContext>-Send>d__7.MoveNext() --- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée --- à System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw() à GreenPipes.Agents.PipeContextSupervisor1.d__7.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Transports.HostConfigurationRetryExtensions.d__0.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Context.BaseConsumeContext.d__611.MoveNext() --- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée --- à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task) à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task) à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult() à MassTransit.Context.BaseConsumeContext.<NotifyFaulted>d__561.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Pipeline.Filters.ConsumerMessageFilter2.<GreenPipes-IFilter<MassTransit-ConsumeContext<TMessage>>-Send>d__4.MoveNext() --- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée --- à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task) à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task) à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult() à GreenPipes.Filters.TeeFilter1.<>c__DisplayClass5_0.<g__SendAsync|1>d.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à GreenPipes.Filters.OutputPipeFilter2.<SendToOutput>d__9.MoveNext() --- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée --- à System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw() à GreenPipes.Filters.OutputPipeFilter2.d__9.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Pipeline.Filters.DeserializeFilter.d__4.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à GreenPipes.Filters.RescueFilter2.<GreenPipes-IFilter<TContext>-Send>d__5.MoveNext() --- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée --- à System.Runtime.ExceptionServices.ExceptionDispatchInfo.Throw() à GreenPipes.Filters.RescueFilter2.d__5.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Pipeline.Filters.DeadLetterFilter.d__3.MoveNext()
--- Fin de la trace de la pile à partir de l'emplacement précédent au niveau duquel l'exception a été levée ---
à System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task)
à System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task)
à System.Runtime.CompilerServices.ConfiguredTaskAwaitable.ConfiguredTaskAwaiter.GetResult()
à MassTransit.Transports.ReceivePipeDispatcher.d__17.MoveNext()

补充异常信息

  • 重试计数在未进入消费者Consume方法时就已增加
  • 出现Azure.RequestFailedException,异常内容:

Service request failed. Status: 401 (Unauthorized)
Content:

return (<Code>401</Code>)
Authorization failed for specified action: Manage,EntityWrite.
TrackingId:1d71f7fe-1627-4fe4-86fb-67446207d0a2_G11, SystemTracker:sb-staging-spoke.servicebus.windows.net:MassTransit_ReceiveFault, Timestamp:2022-12-09T19:38:52
Headers:
Transfer-Encoding: chunked
Strict-Transport-Security: REDACTED
Content-Type: application/xml; charset=utf-8
Date: Fri, 09 Dec 2022 19:38:52 GMT
Server: Microsoft-HTTPAPI/2.0


问题分析与解决

1. 修复权限不足问题(核心原因)

401错误明确提示服务主体没有Manage,EntityWrite权限,这是导致死信操作失败、消费者崩溃的根本原因。

  • 给当前使用的服务主体分配Azure Service Bus Data Owner角色,或创建包含Manage和EntityWrite权限的自定义角色,确保能正常操作死信队列和故障实体。

2. 修正配置中的错误

原配置存在两处明显问题:

  • se.Consumer<MyMessage>()指定的是消息类型而非消费者类型,应改为se.Consumer<MyMessageConsumer>()
  • 全局重试策略会作用于死信处理环节,导致重试计数异常增加,应将重试策略限定在消费者范围内

调整后的订阅端点配置:

cfg.SubscriptionEndpoint<T>(_subscriptionName.Value, se =>
{
    se.Rule = new CreateRuleOptions("Receiver", new SqlRuleFilter("receiver='all' OR receiver='MySpecificReceiver'"));
    
    // 仅对消费者逻辑应用重试
    se.UseMessageRetry(c => c.Immediate(5));
    
    // 指定正确的消费者类型
    se.Consumer<MyMessageConsumer>();
    
    se.ConfigureDeadLetterQueueDeadLetterTransport();
    se.ConfigureDeadLetterQueueErrorTransport();
});

3. 添加故障救援逻辑避免消费者崩溃

添加救援过滤器,捕获处理过程中的异常,避免单个消息失败导致整个接收管道崩溃:

cfg.SubscriptionEndpoint<T>(_subscriptionName.Value, se =>
{
    // ...其他配置
    
    se.UseRescue(r =>
    {
        r.Handle<Exception>();
        r.Send(ctx => ctx.DeadLetter("Rescued failure", ctx.Exception.Message));
    });
});

内容的提问来源于stack exchange,提问作者Noone

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 13:05:18