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

Azure上NServiceBus发布订阅异常:发布者需对应消息处理器问题排查

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:10:02