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

NServiceBus与SQL Server发布订阅异常:消息跨队列重复投递

问题:NServiceBus + SQL Server发布订阅异常——所有消息流入所有队列

问题现象

  • 基于NServiceBus和SQL Server搭建发布订阅项目,包含两个IEvent实现类:OpportunityMessage、PageVisitMessage
  • 对应专属处理器OpportunityMessageHandler、VisitorMessageHandler,分别监听messaging_opportunity_in、messaging_visitors_in队列
  • 发布端消息可正常发送,但触发任意消息都会同时进入两个队列
  • 检查Subscriptions表发现:两个端点均被注册为订阅所有消息类型,与代码预期的专属订阅不符,导致消息重复处理

相关代码

消息类代码

public class OpportunityMessage : IEvent
{
    /// <summary>
    /// Gets or sets the name of the form
    /// </summary>
    public string FormName { get; set; }

    /// <summary>
    /// Gets or sets the page identifier
    /// </summary>
    public byte[] PageIdentifier { get; set; }

    /// <summary>
    /// Gets or sets any supplemental data with a form
    /// </summary>
    public string SupplementalData { get; set; }
}

public class PageVisitMessage : IEvent
{
    /// <summary>
    /// Gets or sets the IP address of the requestor
    /// </summary>
    public string IpAddress { get; set; }

    /// <summary>
    /// Gets or sets the URL that generated a page visit
    /// </summary>
    public string NavigateUrl { get; set; }

    /// <summary>
    /// Gets or sets the referrer
    /// </summary>
    public string Referrer { get; set; }

    /// <summary>
    /// Gets or sets the user agent
    /// </summary>
    public string UserAgent { get; set; }
}

处理器类代码

public class OpportunityMessageHandler : BaseMessageHandler, IHandleMessages<OpportunityMessage>
{   
    // 配置文件中对应值为"messaging_opportunity_in"
    public override string QueueName => this.Configuration["Messaging:OpportunityPublishQueue"];

    public Task Handle(OpportunityMessage message, IMessageHandlerContext context)
    {
        // 消息处理逻辑
        return Task.CompletedTask;
    }
}

public class VisitorMessageHandler : BaseMessageHandler, IHandleMessages<PageVisitMessage>
{   
    // 配置文件中对应值为"messaging_visitors_in"
    public override string QueueName => this.Configuration["Messaging:VisitorPublishQueue"];

    // 注意:此处参数类型错误,应为PageVisitMessage
    public Task Handle(OpportunityMessage message, IMessageHandlerContext context)
    {
        // 消息处理逻辑
        return Task.CompletedTask;
    }
}

BaseMessageHandler初始化代码

/// <summary>
/// Initializes this message handler.
/// </summary>
/// <param name="configuration">The <see cref="IConfiguration"/> instance containing the application's configuration.</param>
public async Task Initialize(IConfiguration configuration)
{
    this.Configuration = configuration;

    // 初始化端点配置
    var listener = new EndpointConfiguration(this.QueueName);
    listener.EnableInstallers();
    listener.SendFailedMessagesTo($"{this.QueueName}_errors");

    // 配置SQL Server传输
    var transport = new SqlServerTransport(Configuration.GetConnectionString("CS"))
    {
        DefaultSchema = this.DefaultSchemaName
    };

    transport.SchemaAndCatalog.UseSchemaForQueue($"{this.QueueName}_errors", this.DefaultSchemaName);
    listener.UseTransport(transport);

    // 配置订阅
    transport.Subscriptions.DisableCaching = true;
    transport.Subscriptions.SubscriptionTableName = new NServiceBus.Transport.SqlServer.SubscriptionTableName(
        this.Configuration["Messaging:Subscriptions"], 
        schema: this.DefaultSchemaName);

    // 启动端点
    this.EndpointInstance = await Endpoint.Start(listener).ConfigureAwait(false);
}

异常原因分析

1. 处理器方法参数类型不匹配

VisitorMessageHandler实现了IHandleMessages<PageVisitMessage>接口,但Handle方法的参数却是OpportunityMessage,这会导致NServiceBus的消息处理器扫描逻辑混乱:

  • 接口声明表明该处理器应处理PageVisitMessage,但方法参数指向OpportunityMessage
  • NServiceBus可能因此错误识别该端点需要订阅两种消息类型,甚至退化为订阅所有IEvent类型

2. 端点未限制消息扫描范围

默认情况下,NServiceBus会扫描当前程序集内所有IHandleMessages<T>实现类。如果两个处理器在同一程序集,且端点初始化时未指定扫描范围,会导致:

  • 每个端点启动时都扫描到两个处理器
  • 每个端点自动订阅两种消息类型,最终表现为订阅所有消息

3. 订阅表残留错误记录

即使修复代码,若Subscriptions表中已有错误的订阅记录(两个端点订阅所有消息),且缓存未及时失效,端点启动时仍会读取旧记录,导致异常持续。

修复步骤

1. 修正处理器方法参数

将VisitorMessageHandler的Handle方法参数改为PageVisitMessage,与接口声明一致:

public class VisitorMessageHandler : BaseMessageHandler, IHandleMessages<PageVisitMessage>
{   
    public override string QueueName => this.Configuration["Messaging:VisitorPublishQueue"];

    // 修正参数类型为PageVisitMessage
    public Task Handle(PageVisitMessage message, IMessageHandlerContext context)
    {
        // 消息处理逻辑
        return Task.CompletedTask;
    }
}

2. 限制端点扫描范围

在BaseMessageHandler的Initialize方法中,为每个端点指定仅扫描自身相关的类型,避免扫描到其他处理器:

public async Task Initialize(IConfiguration configuration)
{
    this.Configuration = configuration;

    var listener = new EndpointConfiguration(this.QueueName);
    listener.EnableInstallers();
    listener.SendFailedMessagesTo($"{this.QueueName}_errors");

    // 添加:限制扫描当前处理器及对应消息类型
    var handlerType = this.GetType();
    var messageType = handlerType.GetInterfaces()
        .Where(i => i.IsGenericType && i.GetGenericTypeDefinition() == typeof(IHandleMessages<>))
        .Select(i => i.GetGenericArguments()[0])
        .FirstOrDefault();
    listener.TypesToScan(new[] { handlerType, messageType });

    var transport = new SqlServerTransport(Configuration.GetConnectionString("CS"))
    {
        DefaultSchema = this.DefaultSchemaName
    };

    transport.SchemaAndCatalog.UseSchemaForQueue($"{this.QueueName}_errors", this.DefaultSchemaName);
    listener.UseTransport(transport);

    transport.Subscriptions.DisableCaching = true;
    transport.Subscriptions.SubscriptionTableName = new NServiceBus.Transport.SqlServer.SubscriptionTableName(
        this.Configuration["Messaging:Subscriptions"], 
        schema: this.DefaultSchemaName);

    this.EndpointInstance = await Endpoint.Start(listener).ConfigureAwait(false);
}

3. 清理订阅表错误记录

手动删除Subscriptions表中属于这两个端点的所有订阅条目,然后重新启动两个端点,让NServiceBus生成正确的订阅记录。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 19:24:57