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
相关产品推荐
相关产品推荐

