MassTransit Azure Service Bus:主题/订阅故障消息消费配置问题咨询
问题一:未创建故障队列/订阅的原因及解决步骤
你当前的配置缺少为TestFaultConsumer配置接收端点的步骤,MassTransit不会自动为未配置端点的消费者创建Azure Service Bus资源。默认情况下,当重试耗尽后MassTransit会自动发布Fault<MyEvent>消息,只需补全端点配置即可。
方案1:手动为故障消费者配置订阅端点
修改MassTransit注册代码,在UsingAzureServiceBus块中添加TestFaultConsumer的订阅端点配置:
services.AddMassTransit(busRegistrationConfigurator => { busRegistrationConfigurator.AddConsumer<TestConsumer, TestConsumerDefinition>(); busRegistrationConfigurator.AddConsumer<TestFaultConsumer>(); busRegistrationConfigurator.UsingAzureServiceBus((busRegistrationContext, serviceBusFactoryConfigurator) => { serviceBusFactoryConfigurator.Host(configuration["ServiceBusSettings:ConnectionString"]); // 配置TestConsumer的订阅端点 serviceBusFactoryConfigurator.SubscriptionEndpoint("test-consumer", "my-topic", cfg => { cfg.ConfigureConsumer<TestConsumer>(busRegistrationContext); }); // 配置TestFaultConsumer的订阅端点,监听Fault<MyEvent>消息 serviceBusFactoryConfigurator.SubscriptionEndpoint<Fault<MyEvent>>("test-fault-consumer", cfg => { cfg.ConfigureConsumer<TestFaultConsumer>(busRegistrationContext); }); }); });
方案2:使用ConfigureEndpoints自动配置所有消费者端点
如果不需要手动指定TestConsumer的主题订阅,可通过ConfigureEndpoints自动为所有消费者创建端点(包括故障消费者):
services.AddMassTransit(busRegistrationConfigurator => { busRegistrationConfigurator.AddConsumer<TestConsumer, TestConsumerDefinition>(); busRegistrationConfigurator.AddConsumer<TestFaultConsumer>(); busRegistrationConfigurator.UsingAzureServiceBus((busRegistrationContext, serviceBusFactoryConfigurator) => { serviceBusFactoryConfigurator.Host(configuration["ServiceBusSettings:ConnectionString"]); // 自动配置所有消费者的端点,包括TestConsumer和TestFaultConsumer serviceBusFactoryConfigurator.ConfigureEndpoints(busRegistrationContext); }); });
注:使用ConfigureEndpoints时,TestConsumer会自动订阅MyEvent全限定类名对应的主题,若需指定自定义主题my-topic,仍需手动配置SubscriptionEndpoint。
问题二:Azure Functions中实现Fault消息消费
完全可以在Azure Functions中实现Fault<MyEvent>消息的消费,以下是两种常用方式:
方式1:使用MassTransit Azure Functions集成
通过MassTransit官方集成,直接用IConsumer<Fault<MyEvent>>接口处理消息:
- 安装NuGet包:
MassTransit.Azure.Functions.ServiceBus - 创建消费者类并配置Startup:
// 故障消费者Function public class TestFaultConsumerFunction : IConsumer<Fault<MyEvent>> { public async Task Consume(ConsumeContext<Fault<MyEvent>> context) { // 编写Fault消息处理逻辑,比如记录日志、发送告警等 var originalMessage = context.Message.Message; var exceptions = context.Message.Exceptions; } } // Functions Startup配置 [assembly: FunctionsStartup(typeof(Startup))] public class Startup : FunctionsStartup { public override void Configure(IFunctionsHostBuilder builder) { builder.Services.AddMassTransitForAzureFunctions(cfg => { cfg.AddConsumer<TestFaultConsumerFunction>(); }); } }
方式2:直接使用ServiceBusTrigger手动处理
若不想依赖MassTransit集成,可使用Azure Functions原生ServiceBusTrigger监听对应主题/订阅:
[FunctionName("TestFaultConsumer")] public async Task Run( [ServiceBusTrigger("Fault-MyNamespace.MyEvent", "test-fault-subscription", Connection = "ServiceBusConnection")] string messageJson, ILogger log) { // 反序列化Fault消息(需引用Newtonsoft.Json或System.Text.Json) var fault = JsonConvert.DeserializeObject<Fault<MyEvent>>(messageJson); // 处理逻辑 log.LogInformation($"处理故障消息,原消息ID:{fault.MessageId}"); }
注:主题名格式为Fault-<MyEvent的全限定类名>(例如MyEvent的命名空间为MyApp.Events,则主题名为Fault-MyApp.Events.MyEvent),这是MassTransit发布Fault消息的默认命名规则。
内容的提问来源于stack exchange,提问作者Stix

