运行Pika RabbitMQ Worker模式时遇ConnectionResetError(10054)错误求助
RabbitMQ Worker模式连接异常问题解决
场景说明
我正在运行一个Worker模式:通过sender.py程序将YouTube URL和索引发送至名为process_queue的队列;receiver.py作为Worker程序消费该队列的消息,运行模型后将JSON输出写入product_queue队列。
sender.py代码
import pika import json import time messages=[ { "url":'', "id": 1 }, { "url": '', "id": 2 } ] connection=pika.BlockingConnection(pika.ConnectionParameters(host='localhost')) channel=connection.channel() channel.queue_declare(queue='process_queue', durable=True) channel.basic_qos(prefetch_count=1) for message in messages: messageJson=json.dumps(message) channel.basic_publish( exchange='', routing_key='process_queue', body=messageJson, ) print(f'Sent: {message["id"]}') connection.close()
receiver.py代码
import pika import sys import os import json from model import youtubeProducts def main(): connection = pika.BlockingConnection(pika.ConnectionParameters('localhost', heartbeat=60)) channel = connection.channel() channel.queue_declare(queue='process_queue',durable=True) channel.queue_declare(queue='product_queue',durable=True) def callback(ch, method, properties, body): message = json.loads(body) youtube_url = message['url'] index=message['id'] print(f'Recieved: {index}') prodList=youtubeProducts(youtube_url, index).run() ch.basic_publish(exchange='', routing_key='product_queue', body=prodList) ch.basic_ack(delivery_tag=method.delivery_tag) channel.basic_qos(prefetch_count=1) channel.basic_consume(queue='process_queue', on_message_callback=callback) channel.start_consuming() if __name__ == '__main__': try: main() except KeyboardInterrupt: print('Interrupted') try: sys.exit(0) except SystemExit: os._exit(0)
报错信息
Traceback (most recent call last): File "reciever.py", line 34, in <module> main() File "reciever.py", line 30, in main channel.start_consuming() File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 1883, in start_consuming self._process_data_events(time_limit=None) File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 2044, in _process_data_events self.connection.process_data_events(time_limit=time_limit) File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 851, in process_data_events self._dispatch_channel_events() File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 567, in _dispatch_channel_events impl_channel._get_cookie()._dispatch_events() File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 1510, in _dispatch_events consumer_info.on_message_callback(self, evt.method, File "reciever.py", line 23, in callback ch.basic_publish(exchange='', routing_key='product_queue', body=prodList) File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 2265, in basic_publish self._flush_output() File "C:\Python\Python38\lib\site-packages\pika\adapters\blocking_connection.py", line 1353, in _flush_output self._connection._flush_output(lambda: self.is_closed, *waiters) File "C:\Python\Python38\lib\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(10054, 'An existing connection was forcibly closed by the remote host', None, 10054, None)
问题描述
运行时出现上述错误,问题看似源于channel.start_consuming()。模型会返回JSON对象,若将其替换为print语句可正常输出正确值,但没有任何内容被写入product_queue队列。
解决步骤
1. 修复消息序列化问题(核心原因)
从报错栈可以看出,问题出在ch.basic_publish调用时连接被重置——pika要求消息体必须是字节串或字符串,你直接传入了模型返回的JSON对象,导致序列化失败,触发RabbitMQ断开连接。
修改receiver.py回调函数中的发布代码,将JSON对象转为JSON字符串:
# 替换原来的basic_publish行 ch.basic_publish(exchange='', routing_key='product_queue', body=json.dumps(prodList))
2. 增强连接健壮性
- 给连接添加重试机制和超时设置,避免单次连接失败导致程序退出
- 给消息添加持久化属性,确保消息不会因RabbitMQ重启丢失
修改receiver.py的连接参数和发布代码:
# 修改连接参数部分 connection = pika.BlockingConnection(pika.ConnectionParameters( host='localhost', heartbeat=60, connection_attempts=3, # 重试连接次数 retry_delay=5 # 重试间隔(秒) )) # 修改消息发布部分 ch.basic_publish( exchange='', routing_key='product_queue', body=json.dumps(prodList), properties=pika.BasicProperties( delivery_mode=2, # 标记消息为持久化 ) )
3. 排查模型运行耗时
如果模型运行时间超过RabbitMQ的心跳间隔(当前设置60秒),会被判定为连接失效。可以:
- 延长心跳时间,比如设置
heartbeat=300 - 若模型运行确实超长时间,考虑拆分任务或使用异步连接(如pika的异步适配器)
4. 验证基础状态
- 检查本地RabbitMQ服务是否正常运行,重启服务尝试
- 通过RabbitMQ管理界面(默认
http://localhost:15672)查看两个队列的状态,确认队列存在且无未确认消息
内容的提问来源于stack exchange,提问作者KeTroy
相关产品推荐
相关产品推荐

