You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

运行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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 12:53:23