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

RabbitMQ升级后.NET Core客户端高负载下超时问题求助

RabbitMQ 版本升级后并发长轮询消费者超时问题

原系统使用RabbitMQ 3.8.5 + RabbitMQ.Client 5.2.0,客户端发起长轮询Web请求时,每个客户端对应创建专属消费者,等待发送至专属队列的特定指令;超时后销毁消费者,客户端重启长轮询,此前可支持1500+并发客户端/消费者。

升级至RabbitMQ 3.13.3 + RabbitMQ.Client 6.8.1后,仅约700个客户端发起长轮询时就触发System.TimeoutException;500个客户端连接超过1小时也会出现相同问题。使用系统其他模块已在用的EasyNetQ消费时,同样报错:超时后RabbitMQ服务器主动终止连接。

进一步排查发现,超时发生时伴随心跳丢失,RabbitMQ服务器终止连接引发连锁故障,目前未找到超时根源及RabbitMQ操作阻塞原因。

补充信息

  • 版本升级前后配置无变更
  • 目标负载仍为1500+客户端
  • 运行环境:Windows Server 2022 Standard x64虚拟机,32GB内存,2.3GHz CPU
  • 队列类型保持经典队列

异常信息

初始异常

at RabbitMQ.Client.Impl.SimpleBlockingRpcContinuation.GetReply(TimeSpan timeout)
at RabbitMQ.Client.Impl.ModelBase.QueueDeclare(String queue, Boolean passive, Boolean 
 durable, Boolean exclusive, Boolean autoDelete, IDictionary`2 arguments)
at RabbitMQ.Client.Impl.ModelBase.ConsumerCount(String queue)

EasyNetQ错误

System.TimeoutException: The operation has timed out.
   at RabbitMQ.Client.Impl.SimpleBlockingRpcContinuation.GetReply(TimeSpan timeout)
   at RabbitMQ.Client.Impl.ModelBase.QueueDeclare(String queue, Boolean passive, Boolean durable, Boolean exclusive, Boolean autoDelete, IDictionary`2 arguments)
   at RabbitMQ.Client.Impl.AutorecoveringModel.QueueDeclare(String queue, Boolean durable, Boolean exclusive, Boolean autoDelete, IDictionary`2 arguments)
   at EasyNetQ.RabbitAdvancedBus.<>c__DisplayClass51_0.<QueueDeclareAsync>b__0(IModel x)
   at EasyNetQ.Persistent.PersistentChannel.InvokeChannelActionAsync[TResult,TChannelAction](TChannelAction channelAction, CancellationToken cancellationToken)
   at EasyNetQ.RabbitAdvancedBus.QueueDeclareAsync(String name, Action`1 configure, CancellationToken cancellationToken)
   at EasyNetQ.AdvancedBusExtensions.QueueDeclare(IAdvancedBus bus, String name, Boolean durable, Boolean exclusive, Boolean autoDelete, CancellationToken cancellationToken)

相关代码

初始消费者实现

ncChannel.QueueDeclare(queueName, true, false, false, null);
ncChannel.ExchangeDeclare(ExchangeName, ExchangeType.Topic, true, false, null);
ncChannel.QueueBind(queueName, ExchangeName, routingKey);
Channel<BasicDeliverEventArgs> responseChannel = Channel.CreateBounded<BasicDeliverEventArgs>(channelOptions);
EventHandler<BasicDeliverEventArgs> nextCommandHandler = async (object sender, BasicDeliverEventArgs msg) =>
{
    try
    {
        if (msg != null && msg.Body.ToArray() != null)
        {
            string msgBody = Encoding.UTF8.GetString(msg.Body.ToArray());
            //log
        }
        else
        {
            //log
        }
        await responseChannel.Writer.WriteAsync(msg);
        responseChannel.Writer.TryComplete();
    }
    catch (Exception ex)
    {
        //log
    }
};
consumer.Received += nextCommandHandler;
 bool autoAck = true;
 consumerTag = ncChannel.BasicConsume(queueName, autoAck, consumer);

异步事件消费者改造尝试

Channel<string> responseChannel = Channel.CreateBounded<string>(channelOptions);

ncChannel.QueueDeclare(queueName, true, false, false, null);
ncChannel.ExchangeDeclare(ExchangeName, ExchangeType.Topic, true, false, null);
ncChannel.QueueBind(queueName, ExchangeName, routingKey);
AsyncEventingBasicConsumer consumer = new AsyncEventingBasicConsumer(ncChannel);
AsyncEventHandler<BasicDeliverEventArgs> nextCommandHandler = async (object sender, BasicDeliverEventArgs msg) =>
{
    try
    {
        if (msg != null && msg.Body.ToArray() != null)
        {
            string msgBody = Encoding.UTF8.GetString(msg.Body.ToArray());
            //log
        }
        else
        {
            //log
        }
        await responseChannel.Writer.WriteAsync(msg);
        responseChannel.Writer.TryComplete();
        await Task.Yield();
    }
    catch (Exception ex)
    {
        //log
    }
};
consumer.Received += nextCommandHandler;

bool autoAck = true;
consumerTag = ncChannel.BasicConsume(queueName, autoAck, consumer);

EasyNetQ消费者实现

var queue = _rabbitBus.Advanced.QueueDeclare(queueName, true, false, false);
consumer = _rabbitBus.Advanced.Consume(queue, (body, properties, info) => Task.Factory.StartNew(async () =>
{
    try
    {
        string msgBody = null;
        if (body.ToArray() != null)
        {
            msgBody = Encoding.UTF8.GetString(body.ToArray());
            //log
        }
        else
        {
            //log
        }
        await responseChannel.Writer.WriteAsync(msgBody);
        responseChannel.Writer.TryComplete();
        await Task.Yield();
    }
    catch (Exception ex)
    {
        //log
    }
}));

内容的提问来源于stack exchange,提问作者Kaizer69

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:07:05