You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.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的参数完全一致,参数不匹配会导致消费者连接到新建的空队列,而非目标队列。
  • 消息路由验证:如果使用交换机,检查生产者是否将消息正确绑定到目标队列;如果用默认交换机,确认发送消息时指定的队列名称和消费者的队列名称完全一致(注意大小写)。
  • 异步消费的资源生命周期问题:
    1. ConsumeMessages中用using var channel,方法执行完毕后通道会被释放,而Received事件是异步触发的,可能在通道释放后才收到消息,导致无法捕获。建议不要在该方法内释放通道,或者改用同步消费方式。
    2. _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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 02:49:55