NServiceBus基于RabbitMQ部署时消息接收失败问题排查
NServiceBus + RabbitMQ 接收服务无法显示且无法接收事件问题
我有一个后端服务,用于发布简单的字符串对象事件。事件发布服务(SimpleSend)在RabbitMQ控制台中可见,但接收/处理服务(SimpleReceive)未显示,且处理程序无法接收到事件。
事件定义
public class TransTest : IEvent { public string TimeSlots { get; set; } public TransTest(string timeSlots) { TimeSlots = timeSlots; } }
发送端(SimpleSender)依赖注入配置(已从.NET 5升级至.NET 6)
本地与托管RabbitMQ环境均已测试:
services.AddSingleton<IMessageSession>(provider => { var endpointConfiguration = new EndpointConfiguration("SimpleSender"); endpointConfiguration.UseSerialization<SystemJsonSerializer>(); endpointConfiguration.AutoSubscribe(); var transport = endpointConfiguration.UseTransport<RabbitMQTransport>(); //transport.ConnectionString("host=localhost"); transport.ConnectionString("host=rabbitmq;username=dev;password=dev"); transport.UseConventionalRoutingTopology(QueueType.Quorum); //transport.UseDirectRoutingTopology( // QueueType.Classic, // exchangeNameConvention: () => "name_event_bus" //); transport.Routing().RouteToEndpoint(typeof(TransTest), "SimpleReceiver"); endpointConfiguration.EnableInstallers(); var endpointInstance = NServiceBus.Endpoint.Start(endpointConfiguration).GetAwaiter().GetResult(); return endpointInstance; });
发布事件的API端点
[HttpGet("/Test")] public async Task<ActionResult> GetTestTrans() { var sendOptions = new SendOptions(); sendOptions.SetDestination("SimpleReceiver"); Console.WriteLine("Sending message to calendar"); var saga = new TransTest("MESSAGE"); await _messageSession.Publish(saga); return Ok(); }
接收端(SimpleReceiver)依赖注入配置
services.AddSingleton<IMessageSession>(provider => { var endpointConfiguration = new EndpointConfiguration("SimpleReceiver"); endpointConfiguration.UseSerialization<SystemJsonSerializer>(); endpointConfiguration.AutoSubscribe(); var transport = endpointConfiguration.UseTransport<RabbitMQTransport>(); //transport.ConnectionString("host=localhost"); transport.ConnectionString("host=rabbitmq;username:dev;password:dev"); transport.UseConventionalRoutingTopology(QueueType.Quorum); //transport.UseDirectRoutingTopology( // QueueType.Classic, // exchangeNameConvention: () => "name_event_bus" // ); endpointConfiguration.EnableInstallers(); var endpointInstance = NServiceBus.Endpoint.Start(endpointConfiguration).GetAwaiter().GetResult(); return endpointInstance; });
事件处理程序
public class SagaTimeSlotsHandler : IHandleMessages<TransTest> { static ILog log = LogManager.GetLogger<SagaTimeSlotsHandler>(); public async Task Handle(TransTest message, IMessageHandlerContext context) { Console.WriteLine("SAGA HANDLED"); Console.WriteLine(message.TimeSlots); log.Info($"Hello from {nameof(SagaTimeSlotsHandler)}"); Task.Completed; } }
RabbitMQ控制台截图

内容的提问来源于stack exchange,提问作者lp_nave
相关产品推荐
相关产品推荐

