Kombu连接RabbitMQ时Pytest测试出现[Errno 104]连接重置问题
Kombu连接RabbitMQ单元测试异常问题
问题场景
使用Kombu连接RabbitMQ时,通用生产者函数在多数业务场景运行正常,但部分单元测试失败:
- 打开
with连接块时连接自动关闭 - 手动调用
conn.connect()触发kombu.exceptions.OperationalError: [Errno 104] Connection reset by peer
所有操作基于Docker Compose,测试在容器内执行。虽错误提示网络问题,但其他场景发消息正常,推测异常与测试配置相关。
已完成验证
- 消息序列化逻辑正常
- 连接字符串可正常使用
- Docker容器间网络连通性无问题
- 版本兼容:Kombu 5.2.3、AMQP 5.1.1、RabbitMQ 3.11.8
- 测试过程无并发连接冲突
- 目标队列无积压消息
错误回溯信息
platform/resources.py:39: in publish_event send_amqp_message(event) messagebus/producer.py:16: in send_amqp_message conn.connect() /usr/local/lib/python3.9/dist-packages/kombu/connection.py:275: in connect return self._ensure_connection( /usr/local/lib/python3.9/dist-packages/kombu/connection.py:434: in _ensure_connection return retry_over_time( /usr/local/lib/python3.9/dist-packages/kombu/utils/functional.py:312: in retry_over_time return fun(*args, **kwargs) /usr/local/lib/python3.9/dist-packages/kombu/connection.py:878: in _connection_factory self._connection = self._establish_connection() /usr/local/lib/python3.9/dist-packages/kombu/connection.py:813: in _establish_connection conn = self.transport.establish_connection() /usr/local/lib/python3.9/dist-packages/kombu/transport/pyamqp.py:201: in establish_connection conn.connect() /usr/local/lib/python3.9/dist-packages/amqp/connection.py:329: in connect self.drain_events(timeout=self.connect_timeout) /usr/local/lib/python3.9/dist-packages/amqp/connection.py:525: in drain_events while not self.blocking_read(timeout): /usr/local/lib/python3.9/dist-packages/amqp/connection.py:530: in blocking_read frame = self.transport.read_frame() /usr/local/lib/python3.9/dist-packages/amqp/transport.py:312: in read_frame payload = read(size) /usr/local/lib/python3.9/dist-packages/amqp/transport.py:627: in _read s = recv(n - len(rbuf)) /usr/local/lib/python3.9/dist-packages/httpretty/core.py:697: in recv return self.forward_and_trace('recv', buffersize, *args, **kwargs) _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ self = <httpretty.core.fakesock.socket object at 0x402d1bf760>, function_name = 'recv', a = (1342177288,), kw = {} function = <built-in method recv of socket object at 0x402d1b19a0> callback = <built-in method recv of socket object at 0x402d1b19a0> def forward_and_trace(self, function_name, *a, **kw): if self.truesock and not self.__truesock_is_connected__: self.truesock = self.create_socket() ### self.connect_truesock() if self.__truesock_is_connected__: function = getattr(self.truesock, function_name) if self.is_http: if self.truesock and not self.__truesock_is_connected__: self.truesock = self.create_socket() ### self.connect_truesock() if not self.truesock: raise UnmockedError() callback = getattr(self.truesock, function_name) > return callback(*a, **kw) E ConnectionResetError: [Errno 104] Connection reset by peer /usr/local/lib/python3.9/dist-packages/httpretty/core.py:666: ConnectionResetError
生产者函数代码
def send_amqp_message(message: BaseEvent, queue: str = "eventbus") -> None: with Connection(MESSAGEBUS_CONNECTION_STRING) as conn: if not conn.connected: conn.connect() queue = conn.SimpleQueue(queue) queue.put({"event": message.name, "payload": asdict(message.data)}, serializer="uuid") queue.close()
RabbitMQ调试日志
rabbitmq-1 | 2023-05-22 17:18:53.633771+00:00 [info] <0.2691.0> accepting AMQP connection <0.2691.0> (172.26.0.5:32944 -> 172.26.0.4:5672) rabbitmq-1 | 2023-05-22 17:18:53.633833+00:00 [error] <0.2691.0> closing AMQP connection <0.2691.0> (172.26.0.5:32944 -> 172.26.0.4:5672): rabbitmq-1 | 2023-05-22 17:18:53.633833+00:00 [error] <0.2691.0> {bad_header,<<1,0,0,0,0,0,177,0>>} rabbitmq-1 | 2023-05-22 17:18:55.687146+00:00 [info] <0.2701.0> accepting AMQP connection <0.2701.0> (172.26.0.5:32958 -> 172.26.0.4:5672) rabbitmq-1 | 2023-05-22 17:18:55.693462+00:00 [warning] <0.2701.0> closing AMQP connection <0.2701.0> (172.26.0.5:32958 -> 172.26.0.4:5672): rabbitmq-1 | 2023-05-22 17:18:55.693462+00:00 [warning] <0.2701.0> client unexpectedly closed TCP connection rabbitmq-1 | 2023-05-22 17:18:55.695395+00:00 [info] <0.2706.0> accepting AMQP connection <0.2706.0> (172.26.0.5:32968 -> 172.26.0.4:5672) rabbitmq-1 | 2023-05-22 17:18:55.695585+00:00 [error] <0.2706.0> closing AMQP connection <0.2706.0> (172.26.0.5:32968 -> 172.26.0.4:5672): rabbitmq-1 | 2023-05-22 17:18:55.695585+00:00 [error] <0.2706.0> {bad_header,<<1,0,0,0,0,0,177,0>>} rabbitmq-1 | 2023-05-22 17:18:59.759150+00:00 [info] <0.2718.0> accepting AMQP connection <0.2718.0> (172.26.0.5:45838 -> 172.26.0.4:5672) rabbitmq-1 | 2023-05-22 17:18:59.765966+00:00 [warning] <0.2718.0> closing AMQP connection <0.2718.0> (172.26.0.5:45838 -> 172.26.0.4:5672): rabbitmq-1 | 2023-05-22 17:18:59.765966+00:00 [warning] <0.2718.0> client unexpectedly closed TCP connection rabbitmq-1 | 2023-05-22 17:18:59.767295+00:00 [info] <0.2723.0> accepting AMQP connection <0.2723.0> (172.26.0.5:45842 -> 172.26.0.4:5672) rabbitmq-1 | 2023-05-22 17:18:59.767437+00:00 [error] <0.2723.0> closing AMQP connection <0.2723.0> (172.26.0.5:45842 -> 172.26.0.4:5672): rabbitmq-1 | 2023-05-22 17:18:59.767437+00:00 [error] <0.2723.0> {bad_header,<<1,0,0,0,0,0,177,0>>}
Docker Compose配置
version: '2' services: rabbitmq: image: rabbitmq:3-management-alpine restart: always environment: RABBITMQ_DEFAULT_USER: guest RABBITMQ_DEFAULT_PASS: guest RABBITMQ_LOG_LEVEL: debug ports: - 5672:5672 - 15672:15672 volumes: - ~/.docker-conf/rabbitmq/data/:/var/lib/rabbitmq/ - ~/.docker-conf/rabbitmq/log/:/var/log/rabbitmq redis: image: redis:alpine test_container: command: /src/bin/debug depends_on: - rabbitmq - redis environment: CELERY_URI: amqp://guest:guest@rabbitmq TRUSTED_DOMAINS: 'service;example.com;another-example.com;test' MESSAGEBUS_CONNECTION_STR: amqp://guest:guest@rabbitmq build: context: . dockerfile: Dockerfile ports: - "8000:8000" volumes: - '.:/src'
Pytest.ini配置
[pytest] norecursedirs = ve testpaths = /src/tests junit_family=xunit2 ; addopts=--tb=short -n 5 --dist=loadscope
问题原因分析
- HTTP Mock库干扰AMQP连接:从错误回溯可见,
httpretty库拦截了AMQP的socket连接请求。httpretty默认会mock所有网络socket,包括非HTTP协议的RabbitMQ 5672端口,导致发送给RabbitMQ的请求头不符合AMQP协议规范(对应日志中的bad_header错误),最终被RabbitMQ主动断开连接。 - 生产者函数冗余代码:
with Connection上下文管理器会自动处理连接的建立与关闭,手动调用conn.connect()可能导致连接状态异常,触发不必要的重连逻辑。 - 容器启动顺序问题:
depends_on仅保证容器启动顺序,不确保RabbitMQ服务完全就绪,测试可能在RabbitMQ未初始化完成时发起连接请求。
解决方案
1. 排除httpretty对AMQP连接的拦截
在涉及RabbitMQ的测试用例中,添加httpretty.allow_net_connect()允许真实网络连接,或指定允许连接的主机/端口:
import httpretty def test_amqp_producer(): # 允许连接到RabbitMQ服务 httpretty.allow_net_connect(allowed_networks=["rabbitmq:5672"]) # 执行测试逻辑 send_amqp_message(test_event)
2. 简化生产者函数代码
移除手动连接的冗余代码,依赖上下文管理器自动处理连接:
def send_amqp_message(message: BaseEvent, queue: str = "eventbus") -> None: with Connection(MESSAGEBUS_CONNECTION_STRING) as conn: queue = conn.SimpleQueue(queue) queue.put({"event": message.name, "payload": asdict(message.data)}, serializer="uuid") queue.close()
3. 确保RabbitMQ服务就绪后再执行测试
修改Docker Compose,为RabbitMQ添加健康检查,并让测试容器等待RabbitMQ就绪:
services: rabbitmq: # ... 原有配置 ... healthcheck: test: ["CMD", "rabbitmq-diagnostics", "check_port_connectivity"] interval: 5s timeout: 10s retries: 5 test_container: # ... 原有配置 ... depends_on: rabbitmq: condition: service_healthy redis: condition: service_started
内容的提问来源于stack exchange,提问作者Niels Uitterdijk
相关产品推荐
相关产品推荐

