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()
其他可能性
我尚未探索的其他可能性: