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

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

问题原因分析

  1. HTTP Mock库干扰AMQP连接:从错误回溯可见,httpretty库拦截了AMQP的socket连接请求。httpretty默认会mock所有网络socket,包括非HTTP协议的RabbitMQ 5672端口,导致发送给RabbitMQ的请求头不符合AMQP协议规范(对应日志中的bad_header错误),最终被RabbitMQ主动断开连接。
  2. 生产者函数冗余代码:with Connection上下文管理器会自动处理连接的建立与关闭,手动调用conn.connect()可能导致连接状态异常,触发不必要的重连逻辑。
  3. 容器启动顺序问题: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 08:32:00