如何避免Python中Pika连接RabbitMQ的生产者断开连接?
RabbitMQ Pika生产者空闲后连接断开问题解决
问题描述
使用Python Pika库搭建RabbitMQ生产者时,若生产者长时间空闲(示例中发送第一条消息后等待5分钟),再次发送消息会触发StreamLostError错误,提示连接被重置。尝试调整心跳机制但未成功,且不确定消息发送频率,不想每次发送都开关通道。
错误信息
channel.basic_publish(exchange="", routing_key=queue_name, body="message 2") File "/usr/local/lib/python3.9/site-packages/pika/adapters/blocking_connection.py", line 2265, in basic_publish self._flush_output() File "/usr/local/lib/python3.9/site-packages/pika/adapters/blocking_connection.py", line 1353, in _flush_output self._connection._flush_output(lambda: self.is_closed, *waiters) File "/usr/local/lib/python3.9/site-packages/pika/adapters/blocking_connection.py", line 523, in _flush_output raise self._closed_result.value.error pika.exceptions.StreamLostError: Stream connection lost: ConnectionResetError(104, 'Connection reset by peer')
原因分析
默认情况下,RabbitMQ的心跳间隔为60秒,若客户端在2倍心跳间隔(120秒)内未发送心跳包,RabbitMQ会主动断开连接。示例中time.sleep(300)会阻塞Pika的BlockingConnection,使其无法发送心跳包,超过120秒后连接被RabbitMQ关闭,再次发送消息时就会报错。
解决方法
1. 配置心跳并定期处理连接事件
通过设置合适的心跳间隔,同时在长时间阻塞期间定期处理连接事件,确保心跳包能正常发送:
import os import pika import time # Connection details RABBITMQ_HOST = os.getenv("RABBITMQ_HOST", "localhost") RABBITMQ_PORT = os.getenv("RABBITMQ_PORT", "5672") if __name__ == "__main__": queue_name = "TEST_QUEUE" # 设置心跳间隔为300秒,覆盖默认的60秒,同时设置阻塞超时 parameters = pika.ConnectionParameters( host=RABBITMQ_HOST, port=RABBITMQ_PORT, heartbeat=300, blocked_connection_timeout=300 ) connection = pika.BlockingConnection(parameters) channel = connection.channel() channel.queue_declare(queue_name, durable=True, auto_delete=False) channel.basic_publish(exchange="", routing_key=queue_name, body="message 1") # 分阶段休眠,每30秒处理一次连接事件以发送心跳 total_sleep = 300 interval = 30 while total_sleep > 0: time.sleep(min(interval, total_sleep)) total_sleep -= interval # 处理心跳及其他连接事件,time_limit=0表示非阻塞处理 connection.process_data_events(time_limit=0) channel.basic_publish(exchange="", routing_key=queue_name, body="message 2") connection.close()
2. 实现自动重连机制
针对不确定消息发送频率的场景,封装生产者类,在每次发送前检查连接状态,断开时自动重建:
import os import pika import time class RabbitMQProducer: def __init__(self, host, port, queue_name): self.host = host self.port = port self.queue_name = queue_name self.connection = None self.channel = None self._connect() def _connect(self): # 清理旧连接 if self.connection and not self.connection.is_closed: self.connection.close() parameters = pika.ConnectionParameters( host=self.host, port=self.port, heartbeat=60 ) self.connection = pika.BlockingConnection(parameters) self.channel = self.connection.channel() self.channel.queue_declare(self.queue_name, durable=True, auto_delete=False) def publish(self, message): try: # 检查连接和通道是否活跃 if not self.connection or self.connection.is_closed: self._connect() if not self.channel or self.channel.is_closed: self.channel = self.connection.channel() self.channel.queue_declare(self.queue_name, durable=True, auto_delete=False) self.channel.basic_publish(exchange="", routing_key=self.queue_name, body=message) except pika.exceptions.AMQPConnectionError: # 捕获连接错误,重连后重试发送 self._connect() self.channel.basic_publish(exchange="", routing_key=self.queue_name, body=message) if __name__ == "__main__": RABBITMQ_HOST = os.getenv("RABBITMQ_HOST", "localhost") RABBITMQ_PORT = os.getenv("RABBITMQ_PORT", "5672") producer = RabbitMQProducer(RABBITMQ_HOST, RABBITMQ_PORT, "TEST_QUEUE") producer.publish("message 1") time.sleep(300) producer.publish("message 2")
3. 心跳机制核心逻辑说明
- 心跳是客户端与RabbitMQ维持连接的信号,客户端每隔
heartbeat秒发送一次心跳包 - RabbitMQ会在
2*heartbeat秒未收到心跳时断开连接 BlockingConnection在阻塞操作(如time.sleep)期间无法处理心跳,需主动调用process_data_events触发心跳发送
内容的提问来源于stack exchange,提问作者Anloubie
相关产品推荐
相关产品推荐

