Azure上NServiceBus发布订阅异常:发布者需对应消息处理器问题排查
问题概述
我有两个不同端点的服务,希望仅通过发布订阅(pub/sub)实现消息交互,未来可能有其他服务需处理这些消息。当前已实现:Service 1(客户端)可发布消息,Service 2(服务端)订阅并处理这些消息。
异常现象
当Service 2在消息处理器内发布响应消息时,消息能正常被Service 1接收,但Service 2端出现错误:
Immediate Retry is going to retry message ... because of an exception: System.InvalidOperationException: No handlers could be found for message type: TP.GetDocumentInfoRes at NServiceBus.LoadHandlersConnector.Invoke... at NServiceBus.ProcessingStatisticsBehavior.Invoke... at NServiceBus.TransportReceiveToPhysicalMessageConnector.Invoke... at NServiceBus.Transport.AzureServiceBus.MessagePump.ProcessMessage...
看起来NServiceBus要求发布消息的Service 2必须拥有该消息类型的处理器,添加空处理器后错误消失。
疑问
我对NServiceBus发布订阅的理解有误吗?如何避免该异常及重试?我期望发布的消息能分发给所有订阅者队列,无需发布者拥有对应处理器,实现“即发即忘”。
已做操作
已实现NServiceBus Ping Pong等发布订阅示例,拆分两个独立服务端点时也出现相同行为。
额外信息
错误消息属性显示Service 2似乎在发布的同时向自身回复,属性如下:
replyTo: Service2_dev NServiceBus.ReplyToAddress: Service2_dev NServiceBus.OriginatingEndpoint: Service2_dev NServiceBus.ProcessingEndpoint: Service2_dev NServiceBus.MessageIntent: Publish
代码示例
消息定义
public interface IMessage : IEvent { ... } public class Message : IMessage { ... } public interface IRequest : IMessage {} public class Request : Message, IRequest { ... } public interface IResponse<T> : IMessage where T : IMessage { public T Request {get;} ... } public abstract class Response<T> : Message, IResponse<T> where T : IMessage { public T Request {get; set;} ... } public interface IGetDocumentInfoRequest : IRequest { GetDocumentInfoCriteria Criteria { get; set; } } [DataContract] public class GetDocumentInfoReq : Request, IGetDocumentInfoRequest { [DataMember] public GetDocumentInfoCriteria Criteria { get; set; } ... } public interface IGetDocumentInfoResponse : IResponse<GetDocumentInfoReq> { GetDocumentInfoResult Result { get; set; } } [DataContract] public class GetDocumentInfoRes : Response<GetDocumentInfoReq>, IGetDocumentInfoResponse, IEvent { [DataMember] public GetDocumentInfoResult Result { get; set; } ... }
两个服务的Program主方法
public class Program { public static void Main(string[] args) { HostApplicationBuilder builder = Host.CreateApplicationBuilder(args); ... EndpointConfiguration endpointConfiguration = new EndpointConfiguration($"{builder.Configuration["ServiceName"]}_{builder.Configuration["Instance"]}"); endpointConfiguration.UseTransport(new AzureServiceBusTransport( builder.Configuration.GetConnectionString("AzureServiceBusConnectionString"))); endpointConfiguration.UseSerialization<SystemJsonSerializer>(); endpointConfiguration.EnableInstallers(); builder.UseNServiceBus(endpointConfiguration); ... } }
Service 1(客户端)
在Quartz任务中调用IMessageSession的publish方法:
await MSG.Publish(new GetDocumentInfoReq( Guid.Parse("db94a8a8-7ade-4ae3-81fe-4b4b578f4444"), Guid.NewGuid(), criteria));
处理响应消息:
public class GetDocumentInfoResponseHandler : ResponseMessageHandler<GetDocumentInfoRes>, IHandleMessages<GetDocumentInfoRes> { ... public async override Task BusinessLogicAsync(GetDocumentInfoRes Message, IMessageHandlerContext Context) { DAL.SaveGetDocumentInfoResult(Message.Result); await Task.CompletedTask; } }
Service 2(服务端)
处理请求并发布响应的消息处理器:
public class GetDocumentInformationRequestHandler : RequestMessageHandler<GetDocumentInfoReq>, IHandleMessages<GetDocumentInfoReq> { ... public async override Task BusinessLogicAsync(GetDocumentInfoReq Message, IMessageHandlerContext Context) { await Context.Publish(new GetDocumentInfoRes(Message, await CLN.GetDocumentInfoAsync(Message.Criteria.ToGetDocumentInfoCriteria(loginResult)), this.ProcessingBegan, DateTime.Now)); // 这里发布消息会导致Service 2抛出无订阅者错误 } }
队列与主题
Azure门户中队列显示正确:
Service1_dev Service2_dev Error <-- 无订阅者错误的存放位置
主题与订阅也正常:
bundle-1 service1_dev $default TP.GetDocumentInfoRes service2_dev $default TP.GetDocumentInfoReq
为Service 2添加空响应处理器后,自动添加了对应订阅
解决方案
问题根源
这个异常的核心原因是NServiceBus的自动订阅机制:当你在端点内使用IMessageHandlerContext.Publish发布事件时,框架会默认尝试给当前端点自动订阅该事件类型。但如果当前端点没有该事件的处理器,就会抛出找不到处理器的错误,触发重试。
另外,从你提供的消息属性来看,ReplyToAddress被设置为Service2自身,这是因为你在处理器上下文(IMessageHandlerContext)中发布消息,上下文会继承原消息的ReplyTo属性,导致框架误判需要自身处理这条发布的消息。
解决方法
1. 使用独立的IMessageSession发布消息(推荐)
不要在处理器的上下文(IMessageHandlerContext)中发布响应事件,而是注入独立的IMessageSession实例来发布。这样发布的消息不会继承原消息的上下文属性(比如ReplyToAddress),也不会触发当前端点的自动订阅逻辑。
修改Service 2的处理器代码:
public class GetDocumentInformationRequestHandler : RequestMessageHandler<GetDocumentInfoReq>, IHandleMessages<GetDocumentInfoReq> { private readonly IMessageSession _messageSession; // 构造函数注入IMessageSession public GetDocumentInformationRequestHandler(IMessageSession messageSession) { _messageSession = messageSession; } public async override Task BusinessLogicAsync(GetDocumentInfoReq Message, IMessageHandlerContext Context) { // 使用IMessageSession发布,而非上下文的Publish方法 await _messageSession.Publish(new GetDocumentInfoRes(Message, await CLN.GetDocumentInfoAsync(Message.Criteria.ToGetDocumentInfoCriteria(loginResult)), this.ProcessingBegan, DateTime.Now)); } }
2. 禁用当前端点的自动订阅
如果必须使用上下文发布,可以通过配置禁用当前端点的自动订阅功能。在Service 2的端点配置中添加:
var transport = endpointConfiguration.UseTransport<AzureServiceBusTransport>(); transport.SubscribeOptions().DisableAutoSubscriptions();
注意:禁用自动订阅后,你需要手动管理所有订阅(比如通过Azure门户创建主题订阅,或者使用代码手动订阅),否则端点将无法接收任何事件。
3. 避免消息类型的歧义检查
检查你的消息定义:GetDocumentInfoRes同时实现了IResponse<T>和IEvent,确保NServiceBus能正确识别这是一个事件类型。可以给消息类型添加[Event]特性(如果使用的是NServiceBus的特性标记),或者确保IEvent是NServiceBus的IEvent接口(而非自定义的同名接口)。
验证方法
修改后,Service 2发布GetDocumentInfoRes时:
- 不会再触发自身的处理器查找逻辑
- 消息会正常发送到对应主题,只有订阅了该事件的Service 1会接收处理
- 不再出现重试和错误日志
内容的提问来源于stack exchange,提问作者gsn1074

