【问题标题】:RabbitMQ c# System.IO.EndOfStreamExceptionRabbitMQ c# System.IO.EndOfStreamException
【发布时间】:2013-12-13 21:14:34
【问题描述】:

当消费者阻止接收来自 SharedQueue 的消息时,我收到以下异常:

Unhandled Exception: System.IO.EndOfStreamException: SharedQueue closed
   at RabbitMQ.Util.SharedQueue.EnsureIsOpen()
   at RabbitMQ.Util.SharedQueue.Dequeue()
   at Consumer.Program.Main(String[] args) in c:\Users\pdecker\Documents\Visual
Studio 2012\Projects\RabbitMQTest1\Consumer\Program.cs:line 33

这是抛出异常时正在执行的代码行:

BasicDeliverEventArgs e = (BasicDeliverEventArgs)consumer.Queue.Dequeue();

到目前为止,我已经看到 rabbitMQ 处于非活动状态时发生的异常。我们的应用程序需要让消费者始终保持连接并监听击键。有谁知道这个问题的原因?有谁知道如何从这个问题中恢复?

提前致谢。

【问题讨论】:

  • RMQ 日志最后一行说什么?您在等待控制台中的输入吗? (默认心跳为 5 秒)
  • 是的,一切都在同一台笔记本电脑上运行:RabbitMQ 服务器,一个应用程序是生产者,另一个应用程序是消费者。以下是日志的最后几行:
  • =INFO REPORT==== 13-Dec-2013::10:55:18 === 接受 AMQP 连接 ([FE80::1577:7B91:63E:FD3A ]:49475 -> [FE80::1577:7B91:63E:FD3A]:5672) =警告报告==== 2013 年 12 月 13 日::11:06:09 === 关闭 AMQP 连接 ([FE80::1577:7B91:63E:FD3A]:49475 -> [FE80::1577:7B91:63E:FD3A]:5672):connection_closed_abruptly
  • 对不起,我过早地按了 ENTER。击键是使用 windows 挂钩读取的,因此可以在任何地方输入按键。生产者应用程序和消费者应用程序都是 C# 控制台应用程序。显示传出/传入消息。
  • 我自己也经历过这种行为,目前正在与 RabbitMQ 工程师确认原因。您的消费应用程序和 RabbitMQ 之间是否有任何类型的负载均衡器。如果是这样,我怀疑您的连接正在被服务器关闭,因为内置的心跳功能被负载均衡器中断了。

标签: c# rabbitmq


【解决方案1】:

消费者与渠道绑定:

var consumer = new QueueingBasicConsumer(channel);

因此,如果通道已关闭,则一旦本地 Queue 被清除,消费者将无法获取任何其他事件。

检查要打开的频道

channel.IsOpen == true

并且队列有可用的事件

if( consumer.Queue.Count() > 0 )

调用前:

BasicDeliverEventArgs e = (BasicDeliverEventArgs)consumer.Queue.Dequeue();

更具体地说,我会在致电Dequeue()之前检查以下内容

if( !channel.IsOpen || !connection.IsOpen )
{
    Your_Connection_Channel_Init_Function();
    consumer = new QueueingBasicConsumer(channel);  // consumer is tied to channel
}

if( consumer.Queue.Any() )
    BasicDeliverEventArgs e = (BasicDeliverEventArgs)consumer.Queue.Dequeue();

【讨论】:

  • 为什么自动恢复本身不这样做?
  • 最好只处理异常。可以在事后询问通道/连接状态以获取详细信息。在当前的代码通道/连接中,检查和出队之间可能会死掉,出队仍然可以抛出异常。
【解决方案2】:

别担心,这只是预期的行为,这意味着队列中没有消息需要处理。不试也不行……

consumer.Queue.Any()

只需捕获 EndOfStreamException:

private void ConsumeMessages(string queueName)
{
    using (IConnection conn = factory.CreateConnection())
    {
        using (IModel channel = conn.CreateModel())
        {
            var consumer = new QueueingBasicConsumer(channel);
            channel.BasicConsume(queueName, false, consumer);
            Trace.WriteLine(string.Format("Waiting for messages from: {0}", queueName));

            while (true)
            {
                BasicDeliverEventArgs ea = null;
                try
                {
                    ea = consumer.Queue.Dequeue();
                }
                catch (EndOfStreamException endOfStreamException)
                {
                    Trace.WriteLine(endOfStreamException);
                    // If you want to end listening end of queue call break;
                    break;    
                }
                if (ea == null) break;
                var body = ea.Body;
                // Consume message how you want
                Thread.Sleep(300);
                channel.BasicAck(ea.DeliveryTag, false);
            }
        }
    }
}

【讨论】:

    【解决方案3】:

    还有另一个可能的问题来源:您的公司防火墙。

    那是因为这样的防火墙可以在连接空闲一段时间后断开与 RabbitMQ 的连接。

    虽然 RabbitMQ 连接有一个heartbeat 功能来防止这种情况,但如果在防火墙连接超时之后发生心跳脉冲,它是没有用的。

    这是默认的心跳间隔configuration,以秒为单位:

    默认值:60(3.5.5 版本之前为 580)

    来自RabbitMQ使用心跳检测死 TCP 连接 简介

    网络可能会以多种方式出现故障,有时非常微妙(例如高 丢包率)。中断的 TCP 连接需要相当长的时间 时间(在 Linux 上使用默认配置大约需要 11 分钟,对于 例如)被操作系统检测到。 AMQP 0-9-1 提供了一个 心跳功能,保证应用层及时发现 关于中断的连接(并且完全没有响应 同行)。 心跳还可以防御某些网络设备 可能会终止“空闲”TCP 连接。

    这发生在我们身上,我们通过减少全局配置中的 Heartbeat Timeout Interval 解决了这个问题:

    在您的 rabbitmq.config 中,找到 heartbeat 并将其设置为小于防火墙规则的值。

    你也可以change the interval in your client

    使用 Java 客户端启用心跳 配置心跳 Java 客户端中的超时,将其设置为 ConnectionFactory#setRequestedHeartbeat 创建连接前:

    ConnectionFactory cf = new ConnectionFactory();
    
    // set the heartbeat timeout to 60 seconds
    cf.setRequestedHeartbeat(60);
    

    使用 .NET 客户端启用心跳 配置心跳 .NET 客户端中的超时,将其设置为 ConnectionFactory.RequestedHeartbeat 在创建连接之前:

    var cf = new ConnectionFactory();
    //set the heartbeat timeout to 60 seconds
    cf.RequestedHeartbeat = 60;
    

    【讨论】:

      【解决方案4】:

      这里说这是预期行为的答案是正确的,但是我认为让它通过这样的设计抛出异常是不好的。

      来自the documentation:“如果在其他线程调用 Enqueue() 或队列关闭之前没有可用的项目,则 Dequeue() 的调用者将阻塞。在后一种情况下,此方法将抛出 EndOfStreamException。”

      因此,就像 GlenH7 所说,您必须在调用 Dequeue() (IModel.IsOpen) 之前检查通道是否打开。

      但是,如果通道在 Dequeue() 阻塞时关闭怎么办?我认为最好调用 Queue.DequeueNoWait(null),并通过等待它返回不为 null 的内容来自己阻塞线程。所以,类似:

      while(channel.IsOpen)
      {
          var args = consumer.Queue.DequeueNoWait(null);
          if(args == null) continue;
          //...
      }
      

      这样,它就不会抛出异常。

      【讨论】:

      • 无需其他等待,此代码将热忙循环。队列中什么都没有?好的,让我们马上再次尝试查看队列!例外很好;使用 DequeueNoWait 的唯一真正原因是如果使用与 Rabbit 无关的另一个等待构造(例如 Task.Wait、Handle.Wait)。
      猜你喜欢
      • 1970-01-01
      • 2022-12-19
      • 2017-01-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-08-15
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多