使用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_0
1.<<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:Authorization failed for specified action: Manage,EntityWrite.return (<Code>401</Code>)
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

