【问题标题】:How to reconnect to RabbitMQ?如何重新连接到 RabbitMQ?
【发布时间】:2016-05-13 14:40:16
【问题描述】:

一旦我的 python 脚本从另一个数据源接收到消息,它就必须不断地向 RabbitMQ 发送消息。 python 脚本发送它们的频率可能会有所不同,例如 1 分钟到 30 分钟。

这是我建立与 RabbitMQ 的连接的方法:

  rabt_conn = pika.BlockingConnection(pika.ConnectionParameters("some_host"))
  channel = rbt_conn.channel()

我刚刚遇到异常

pika.exceptions.ConnectionClosed

如何重新连接到它?最好的方法是什么?有什么“策略”吗?是否能够发送 ping 以保持连接处于活动状态或设置超时?

任何指针将不胜感激。

【问题讨论】:

  • 你可以试试 gearman 模块,看看这个gearmanhq.com/help/tutorials/Python/basic。用于创建客户端和工作人员连接。在这里,您的 python 脚本将是客户端,而 rabbitmq 将是工作人员
  • @PrashantPuri,我不想使用任何 API 来完成如此简单的任务!
  • @PrashantPuri,您是否使用Web服务api打开终端?还是重启你的电脑?
  • 它不是 api,它是通用的应用程序框架,允许您并行工作、负载平衡处理以及在语言之间调用函数。
  • @PrashantPuri,你用框架加2和2吗?或者你只是打电话给+?

标签: python rabbitmq pika


【解决方案1】:

RabbitMQ 使用 heartbeats 来检测和关闭“死”连接,并防止网络设备(防火墙等)终止“空闲”连接。从 3.5.5 版开始,默认超时设置为 60 秒(之前约为 10 分钟)。来自docs

心跳帧大约每超时/ 2 秒发送一次。在两次错过的心跳之后,对等方被认为是不可达的。

Pika 的 BlockingConnection 的问题在于它无法响应心跳,直到进行某些 API 调用(例如,channel.basic_publish()connection.sleep() 等)。

到目前为止我发现的方法:

增加或停用超时时间

RabbitMQ 在建立连接时与客户端协商超时。理论上,应该可以使用 heartbeat_interval 参数用更大的默认值覆盖服务器默认值,但当前 Pika 版本 (0.10.0) 使用的 min 值介于服务器和客户端。此问题已在当前 master 上修复。

另一方面,可以通过将heartbeat_interval 参数设置为0 来完全停用心跳功能,这很可能会导致您遇到新问题(防火墙断开连接等)

重新连接

扩展@itsafire 的答案,您可以编写自己的 publisher 类,让您在需要时重新连接。一个简单的实现示例:

import logging
import json
import pika

class Publisher:
    EXCHANGE='my_exchange'
    TYPE='topic'
    ROUTING_KEY = 'some_routing_key'

    def __init__(self, host, virtual_host, username, password):
        self._params = pika.connection.ConnectionParameters(
            host=host,
            virtual_host=virtual_host,
            credentials=pika.credentials.PlainCredentials(username, password))
        self._conn = None
        self._channel = None

    def connect(self):
        if not self._conn or self._conn.is_closed:
            self._conn = pika.BlockingConnection(self._params)
            self._channel = self._conn.channel()
            self._channel.exchange_declare(exchange=self.EXCHANGE,
                                           type=self.TYPE)

    def _publish(self, msg):
        self._channel.basic_publish(exchange=self.EXCHANGE,
                                    routing_key=self.ROUTING_KEY,
                                    body=json.dumps(msg).encode())
        logging.debug('message sent: %s', msg)

    def publish(self, msg):
        """Publish msg, reconnecting if necessary."""

        try:
            self._publish(msg)
        except pika.exceptions.ConnectionClosed:
            logging.debug('reconnecting to queue')
            self.connect()
            self._publish(msg)

    def close(self):
        if self._conn and self._conn.is_open:
            logging.debug('closing queue connection')
            self._conn.close()

其他可能性

我尚未探索的其他可能性:

【讨论】:

    【解决方案2】:

    很简单:像这样的一些模式。

    import time
    
    while True:
        try:
            communication_handles = connect_pika()
            do_your_stuff(communication_handles)
        except pika.exceptions.ConnectionClosed:
            print 'oops. lost connection. trying to reconnect.'
            # avoid rapid reconnection on longer RMQ server outage
            time.sleep(0.5) 
    

    您可能需要重构您的代码,但基本上它是关于捕获异常、缓解问题并继续执行您的工作。 communication_handles 包含所有 pika 元素,例如通道、队列以及您的东西需要通过 pika 与 RabbitMQ 通信的任何内容。

    【讨论】:

    • 很好,但是这个解决方案将进入最大递归深度,以防服务器停止更长时间。
    • 同意。这只是如何处理断开连接的指针。我添加了一个可选建议如何处理 langer 服务器中断。还可以添加一个宽限计数器,以在进行了太多重新连接尝试时最终退出。把自己打晕。
    • 可以与pypi.org/project/retrying结合使用,在放弃前重试N次,以避免无限的重新连接尝试/重试/递归等。
    猜你喜欢
    • 1970-01-01
    • 2014-10-21
    • 1970-01-01
    • 2017-03-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-02-07
    • 1970-01-01
    相关资源
    最近更新 更多