消息代理宕机时,如何让MassTransit抛出异常并手动处理?
我目前用MassTransit搭配RabbitMQ,日志用log4net记录,日常运行都正常。但一旦RabbitMQ服务器宕机,日志里就会刷满这类错误:
ERROR - RabbitMQ Connect Failed: Broker unreachable: localhost:5672/
而且这时候还能无限制发布消息,但这些消息发出去后都会丢失。我想知道能不能触发这类异常并手动处理?或者有没有办法让Publish方法在消息代理宕机时直接抛出异常?
我的总线配置代码如下:
var busControl = Bus.Factory.CreateUsingRabbitMq(cfg => { var host = cfg.Host(new Uri("rabbitmq://localhost/"), h => { }); cfg.UseLog4Net(); cfg.ReceiveEndpoint("test-queue", ep => { ep.StateMachineSaga(context.Resolve<ProductSaga>(), context.Resolve<ILifetimeScope>()); if (ep is IRabbitMqReceiveEndpointConfigurator) { ((IRabbitMqReceiveEndpointConfigurator)ep).PrefetchCount = 8; } ep.UseInMemoryOutbox(); }); });
针对你的问题,我分几个部分来给出具体的解决办法:
一、让Publish在代理不可用时立即抛出异常
MassTransit默认会把无法发送的消息放到内存重试队列中,持续尝试连接代理直到恢复——这就是你能继续发消息但会丢失的核心原因(内存中的消息如果应用重启就彻底没了)。如果想让Publish在代理不可用时直接抛出异常,可以通过以下配置实现:
禁用发送重试
在总线配置里添加发送重试规则,设置为不重试,这样第一次发送失败就会立刻抛出异常:cfg.UseSendRetry(r => r.None());设置发送超时
同时配置发送超时,避免无限等待连接响应:cfg.ConfigurePublish(p => p.UseSendTimeout(TimeSpan.FromSeconds(5)));把这两个配置加到你的总线创建代码中,修改后完整代码如下:
var busControl = Bus.Factory.CreateUsingRabbitMq(cfg => { var host = cfg.Host(new Uri("rabbitmq://localhost/"), h => { }); cfg.UseLog4Net(); // 添加发送重试和超时配置 cfg.UseSendRetry(r => r.None()); cfg.ConfigurePublish(p => p.UseSendTimeout(TimeSpan.FromSeconds(5))); cfg.ReceiveEndpoint("test-queue", ep => { ep.StateMachineSaga(context.Resolve<ProductSaga>(), context.Resolve<ILifetimeScope>()); if (ep is IRabbitMqReceiveEndpointConfigurator) { ((IRabbitMqReceiveEndpointConfigurator)ep).PrefetchCount = 8; } ep.UseInMemoryOutbox(); }); });配置完成后,当RabbitMQ宕机时调用
Publish就会直接抛出异常,你可以在业务代码里捕获这个异常做自定义处理,比如提示用户操作失败、记录专属告警日志,或者把消息暂存到本地持久化存储中。
二、监听连接失败事件做手动处理
如果你不想让Publish直接抛异常,但又想在代理宕机时执行自定义逻辑(比如发送运维告警、暂停业务发布等),可以订阅MassTransit提供的总线连接事件:
监听主机连接失败事件
总线启动后,订阅HostConnectFaulted事件,在连接失败时触发你的处理逻辑:busControl.ConnectHostConnectFaultedHandler(context => { // 这里写你的自定义处理逻辑,比如发送告警邮件、推送通知 Console.WriteLine("RabbitMQ连接失败,已触发自定义告警流程"); return Task.CompletedTask; });监听总线故障事件
还可以订阅更全面的总线故障事件,捕获总线层面的异常:busControl.ConnectBusFaultedHandler(context => { Console.WriteLine($"总线出现故障: {context.Exception.Message}"); // 这里可以加入重启总线、通知运维等逻辑 return Task.CompletedTask; });
三、解决消息丢失问题
另外,针对消息丢失的情况,你可以用持久化Outbox替代内存Outbox,这样当代理宕机时,发布的消息会被持久化到数据库(比如SQL Server、PostgreSQL),等代理恢复后自动重试发送,彻底避免消息丢失。示例配置(以EntityFramework为例,需安装对应NuGet包):
cfg.UseEntityFrameworkOutbox<MyDbContext>(o => { o.QueryDelay = TimeSpan.FromSeconds(10); o.UseSqlServer(); // 根据你的数据库类型选择对应方法 });
配置后即使应用重启,消息也会保存在数据库中,RabbitMQ恢复后会自动完成发送。
内容的提问来源于stack exchange,提问作者Maksim Alkhovik

