.NET 6+SignalR通知服务:RabbitMQ队列长度为0且无法消费消息
问题:基于.NET 6.0的SignalR通知服务无法消费RabbitMQ消息,队列长度显示为0
我正在开发一个基于.NET 6.0的通知服务项目,使用SignalR实现实时功能,计划通过RabbitMQ存储通知队列。目前能够正常发送通知,但无法从RabbitMQ中消费消息,且RabbitMQ队列长度返回0,求排查问题原因。
相关代码
NotificationHub的OnConnectedAsync方法
public override async Task OnConnectedAsync() { try { var token = Context.GetHttpContext().Request.Query["access_token"]; if (!string.IsNullOrEmpty(token)) { var result = await UserConnectionAsync(token); await Clients.Client(Context.ConnectionId).ConnectionResponse(result.response); if (result.succesfulConnection) { await SendStoredNotifications(result.userId); //await _rabbit.ConsumeMessages(result.userId); } else { var response = new ServerResponse { status = "Rejected", message = "Error en la autorizacion" }; await Clients.Client(Context.ConnectionId).ConnectionResponse(response); Context.Abort(); } } else { var response = new ServerResponse { status = "Rejected", message = "Token no proporcionado" }; await Clients.Client(Context.ConnectionId).ConnectionResponse(response); Context.Abort(); } } catch (Exception ex) { Console.WriteLine($"Error {ex.Message}"); throw new Exception($"Error en la conexion {ex.Message}"); } }
SendStoredNotifications方法
private async Task SendStoredNotifications(string userId) { if (!string.IsNullOrEmpty(userId)) { //var notifications = _notificationStore.RetrieveNotifications(userId); var notifications = await _rabbit.ConsumeMessages(userId); foreach (var notification in notifications) { await Clients.Client(Context.ConnectionId).SendNotification(notification); } _notificationStore.ClearNotifications(userId); } }
RabbitMQ连接创建方法
public void CreateConnection() { try { var factory = new ConnectionFactory() { HostName = "localhost", UserName = "guest", Password = "guest", Port = 5672 }; _connection = factory.CreateConnection(); } catch (Exception ex) { throw new BadHttpRequestException(ex.Message); } }
消费消息的ConsumeMessages方法(无法正常工作)
public async Task<List<NotificationModel>> ConsumeMessages(string userId) { List<NotificationModel> messages = new List<NotificationModel>(); var factory = new ConnectionFactory() { HostName = "localhost", UserName = "guest", Password = "guest", Port = 5672 }; CreateConnection(); if (_connection.IsOpen) { using var channel = _connection.CreateModel(); channel.QueueDeclare("Notification", durable: true, exclusive: false, autoDelete: false); var consumer = new EventingBasicConsumer(channel); consumer.Received += (model, ea) => { try { var body = ea.Body.ToArray(); var jsonMessage = Encoding.UTF8.GetString(body); var notification = JsonSerializer.Deserialize<NotificationUserId>(jsonMessage); // Verifica si el ID de usuario coincide con el usuario especificado if (notification.UserId == userId) { messages.Add(notification.Notification); } } catch (Exception ex) { throw new Exception("Erro deserializando el mensaje"); } }; channel.BasicQos(0,prefetchCount: 1, global: false); var result = channel.BasicConsume(queue: "Notification", autoAck: false, consumer: consumer); _messageReceiveEvent.WaitOne(TimeSpan.FromSeconds(10)); channel.ModelShutdown += (model, reason) => { _messageReceiveEvent.Set(); }; } else { throw new Exception("Connection closed"); } return messages; }
目前messages.Count始终为0。
排查方向
- 队列参数一致性:确认生产者创建队列时的参数(
durable、exclusive、autoDelete等)和消费者QueueDeclare的参数完全一致,参数不匹配会导致消费者连接到新建的空队列,而非目标队列。 - 消息路由验证:如果使用交换机,检查生产者是否将消息正确绑定到目标队列;如果用默认交换机,确认发送消息时指定的队列名称和消费者的队列名称完全一致(注意大小写)。
- 异步消费的资源生命周期问题:
ConsumeMessages中用using var channel,方法执行完毕后通道会被释放,而Received事件是异步触发的,可能在通道释放后才收到消息,导致无法捕获。建议不要在该方法内释放通道,或者改用同步消费方式。_messageReceiveEvent.WaitOne(10秒)的逻辑有问题,ModelShutdown事件只有通道关闭时才会触发,10秒内通道没关闭的话,方法会超时返回,此时可能还没收到消息。
- 消息确认机制:设置了
autoAck: false,但Received事件中没有调用channel.BasicAck(ea.DeliveryTag, multiple: false)确认消息,RabbitMQ会认为消息未处理,重新投递但消费者可能无法接收,同时队列中消息会处于未确认状态,影响队列长度统计。 - 用户ID匹配验证:在
Received事件中添加日志,输出接收到的notification.UserId和传入的userId,检查是否存在格式、大小写不匹配的情况,导致消息被过滤。 - 连接管理问题:
ConsumeMessages中重复创建了ConnectionFactory和连接,建议统一管理RabbitMQ连接,避免重复创建导致的资源问题。
内容的提问来源于stack exchange,提问作者Elias
相关产品推荐
相关产品推荐

