【问题标题】:Never ending message loop: Same message redelivered in python rabbitmq consumer永无止境的消息循环:在 python rabbitmq 消费者中重新传递了相同的消息
【发布时间】:2017-04-14 21:09:42
【问题描述】:

我正在使用此处发布的示例消费者:

http://pika.readthedocs.org/en/latest/examples/asynchronous_consumer_example.html

我使用 ExampleConsumer 的原因是当工作任务开始花费更长的时间时,我与 rabbitmq 的连接失败,其中时间超过 10 分钟。在长时间运行的任务完成并且进程失败后,连接被关闭。之前,它需要一分钟左右的时间来处理 1000 条消息。

ExampleConsumer 似乎可以重新连接,但是,在确认消息中,该消息实际上并未得到确认,因为连接已失效。它似乎从下面的确认消息方法正常返回。然后它会尝试重新连接,然后重新发送刚刚完成的消息。

  def acknowledge_message(self, delivery_tag):
        """Acknowledge the message delivery from RabbitMQ by sending a
        Basic.Ack RPC method for the delivery tag.

        :param int delivery_tag: The delivery tag from the Basic.Deliver frame

        """
        LOGGER.info('Acknowledging message %s', delivery_tag)
        self._channel.basic_ack(delivery_tag)

【问题讨论】:

    标签: python rabbitmq pika


    【解决方案1】:

    RabbitMQ 代理实现了一个默认的心跳超时,取决于 RabbitMQ 版本,大约为 10 分钟或 1 分钟;较短的默认值是在较新的版本中,从 RabbitMQ v3.5.5 开始。应用程序可以通过连接参数传递明确的更长的心跳超时首选项。 Pika 的 SelectConnection 没有后台线程,所以当工作任务耗时超过心跳超时时间时,SelectConnection 无法在 broker 预期的时限内服务心跳,broker 会断开连接。您可以通过多种方式尝试解决此问题:

    1. 通过 pika.connection.ConnectionParameters 设置更长的心跳超时首选项(可能是最简单的)。 ConnectionParameters.heartbeat_interval=0 应该完全禁用心跳(和心跳超时)。
    2. 在与任务处理逻辑不同的线程上运行连接
    3. 切换到一种协作式多任务连接类型,例如 Pika 中基于 Tornado 或 Twisted 框架的适配器,或 Haigha 中基于 gevent 的适配器。此更改将要求任务处理逻辑对协作式多任务处理友好。

    【讨论】:

    • 如果您查看上面的异步示例,我正在使用 SelectConnection。我看不到在该连接类型中可配置的心跳。
    • 这个有没有 process_data_events() pika.readthedocs.org/en/latest/modules/adapters/select.html
    • @JustinThomas,SelectConnection 也没有后台线程。见pika.connection.ConnectionParameters。您可以将 ConnectionParameters 传递给 pika 的任何连接。一种选择是将heartbeat_interval=0 传递给ConnectionParameters 构造函数以禁用心跳。
    • 关于heartbeat_interval 的ConnectionParameters 文档字符串中让我有点惊讶的是,它声称将使用服务器和用户提案的最小值,这表明用户无法请求比服务器建议的更大的东西。这似乎不对。 ":param int heartbeat_interval: 发送心跳的频率。将使用此值与服务器提议之间的最小值。使用 0 停用心跳,使用 None 接受服务器提议。"
    • 仍然无法正常工作,尝试显式 heatbeat_interval=0。不确定它应该更高还是更低。它在大约 10 分钟后死亡。尝试将其扩展到 100 个节点并提取我们所有的数据...感谢您的帮助。
    【解决方案2】:

    您可能需要add a heartbeat to your message consumer 才能保持连接有效。

    如果rabbitmq 认为消费者在消息处于“未确认”模式(仍在处理中)时死亡,它会将消息放回队列中。有心跳可能有助于保持连接活跃,防止这种情况发生。

    【讨论】:

    • 我正在使用 SelectConnection 并没有看到任何关于心跳的信息。
    【解决方案3】:

    如果您使用的是 pika 异步消费者示例,则只需将此更改添加到 init 方法:

        self._url = 'amqp://{}:{}@{}:{}/%2F{}'.format(
                         self.USERNAME, self.PASSWORD, self.ADDRESS, self.PORT, self.QUERY) 
    

    使用 self.QUERY 可以参数化设置不同参数的字符串,例如心跳如下:

    self.QUERY ='?heartbeat_interval=600'
    

    connect 方法会处理心跳事务。

    def connect(self):
        """This method connects to RabbitMQ, returning the connection handle.
        When the connection is established, the on_connection_open method
        will be invoked by pika.
    
        :rtype: pika.SelectConnection
    
        """
        LOGGER.info('Connecting to %s', self._url)
        return pika.SelectConnection(parameters=pika.URLParameters(self._url),
                                     on_open_callback=self.on_connection_open,
                                     on_open_error_callback=self.on_connection_error,
                                     stop_ioloop_on_close=False,
                                     )
    

    这是告诉 RabbitMQ 哪个心跳与你的消费者相关联的一个很好的方法。请注意,RabbitMQ 将强制它至少为 60 秒。因此,您不能将其设置得更低。

    有关这些连接参数的更多信息: https://pika.readthedocs.io/en/latest/modules/parameters.html

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-09-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多