使用pika basic_get从RabbitMQ取消息遇StreamLostError求助
使用Pika的basic_get消费RabbitMQ消息时遇StreamLostError问题
我用Pika的basic_get方法从RabbitMQ获取消息,并用basic_ack确认,代码如下:
import pika connection = pika.BlockingConnection( pika.ConnectionParameters( host="rabbitmq", heartbeat=0 ) ) channel = connection.channel() def mq_record(): method_frame, _, body = channel.basic_get("test") if not body: return False, None if method_frame: channel.basic_ack(delivery_tag=method_frame.delivery_tag, multiple=True) rev_mq = eval(body.decode()) task_id = rev_mq.get("task_id") level = rev_mq.get("level") return task_id, level task_id, level = mq_record()
运行后抛出错误:
Traceback (most recent call last): File "/usr/local/lib/python3.7/site-packages/pika/adapters/blocking_connection.py", line 2159, in basic_get self._basic_getempty_result.is_ready) File "/usr/local/lib/python3.7/site-packages/pika/adapters/blocking_connection.py", line 1335, in _flush_output self._connection._flush_output(lambda: self.is_closed, *waiters) File "/usr/local/lib/python3.7/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')
我查到相关说法:
消费进程可能需要太长时间才能完成并向服务器发送Ack/Nack,导致服务器未收到客户端心跳,从而停止服务,客户端因此报错。
但我的场景里确认操作耗时很短,为什么还是出现这个错误?
可能的原因及解决方向
1. 心跳配置完全禁用的问题
你设置了heartbeat=0,这会彻底关闭心跳机制。RabbitMQ服务器会主动断开长时间无交互的连接——哪怕你的确认操作很快,只要两次basic_get的间隔超过服务器默认的connection_timeout(600秒),就会被判定为无效连接并断开。建议恢复合理的心跳值,同时增加连接阻塞超时:
connection = pika.BlockingConnection( pika.ConnectionParameters( host="rabbitmq", heartbeat=60, # 启用心跳,间隔60秒 blocked_connection_timeout=300 # 避免连接被意外阻塞后直接断开 ) )
2. basic_get轮询模式的局限性
basic_get是主动轮询拉取,如果队列长时间为空,客户端会处于无交互的等待状态,服务器可能判定连接失效。可以改用basic_consume的推模式(RabbitMQ主动推送消息,保持连接活跃),或者在轮询间隙加入短sleep,同时依赖Pika的自动心跳维持连接。
3. 网络或服务器层面的问题
ConnectionResetError(104)也可能是网络波动、防火墙拦截、RabbitMQ服务器重启/资源过载导致的:
- 查看RabbitMQ服务器日志,确认是否有主动断开连接的记录
- 测试客户端到服务器的网络稳定性,比如持续ping或telnet检测
- 检查RabbitMQ服务器的CPU、内存、磁盘资源是否充足
4. 消息解析的潜在风险(附加建议)
虽然和当前错误无直接关联,但eval(body.decode())存在严重安全风险,建议改用json.loads解析消息体:
import json # ... rev_mq = json.loads(body.decode())
5. 连接复用问题
如果这段代码被频繁调用,每次都新建连接,可能导致服务器连接数耗尽,从而主动断开连接。建议复用连接和通道,不要每次拉取消息都重新创建。
内容的提问来源于stack exchange,提问作者haojie
相关产品推荐
相关产品推荐

