Python 基础体系 · 第 108/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 集成测试:Testcontainers、数据库、消息队列和稳定隔离

集成测试验证的不是单个函数的返回值,而是多个真实组件连接起来后的行为。例如,应用收到 HTTP 请求后写入 PostgreSQL,再发布一条消息;消费者从 RabbitMQ 或 Kafka 取出消息,执行业务逻辑并更新数据库。这里的“集成”至少包含三类边界:

  1. 进程边界:应用与数据库、消息代理运行在不同进程中。
  2. 协议边界:应用必须遵守 SQL、事务、AMQP、Kafka 协议及客户端库的行为。
  3. 状态边界:数据库记录、队列中的消息、消费者位点和容器文件系统都会在测试过程中发生变化。

单元测试可以用替身对象验证“代码是否调用了某个方法”,但无法证明 SQL 真的能执行、事务真的会提交、消息真的能被代理接受,或者消费者重启后是否会重复处理。因此,集成测试的核心问题不是“如何启动一个容器”,而是:

如何创建一个真实、可观测、可重复、可清理,并且不会被其他测试污染的外部系统。

一、集成测试的边界:真实组件和可控环境

1. 什么是集成测试

设被测应用为 AA,外部依赖为 D1,D2,,DnD_1, D_2, \ldots, D_n。一个集成测试验证的是:

AD1D2A \leftrightarrow D_1 \leftrightarrow D_2 \ldots

在给定初始状态 S0S_0、输入 II 和环境配置 EE 下,系统是否产生预期结果 OO

F(A,D,S0,I,E)=OF(A, D, S_0, I, E) = O

其中:

  • AA:应用代码;
  • DD:真实数据库、消息队列、缓存或第三方协议实现;
  • S0S_0:测试开始前的数据库和消息系统状态;
  • II:HTTP 请求、数据库操作或发布的消息;
  • EE:连接地址、认证信息、超时和功能开关;
  • OO:数据库状态、消息状态、返回结果和副作用。

稳定集成测试至少要求:

S0 可重建S_0 \text{ 可重建}

并且:

F(A,D,S0,I,E)F(A, D, S_0, I, E)

在相同代码、镜像和输入下能够得到相同的可观察结果。

如果测试直接连接开发机上长期运行的数据库,那么 S0S_0 通常不可控:上一个测试可能留下数据,开发者手工操作可能改变数据,数据库版本也可能与 CI 不同。此时测试失败后,很难判断是代码缺陷、环境差异还是脏数据造成的。

2. Testcontainers 解决了什么问题

Testcontainers 是一个在测试运行期间启动 Docker 容器的 Python 库。它提供数据库、Kafka、RabbitMQ、Redis 等常见服务的容器封装,也允许使用能够运行在 Docker 中的自定义服务。当前 Python 实现的项目元数据已声明支持 Python 3.14;其常见用法是通过上下文管理器启动容器、获取连接地址,并在退出上下文时清理资源。(github.com)

它解决的是环境生命周期问题:

测试进程
  │
  ├── 创建容器
  ├── 等待服务真正可用
  ├── 读取动态端口和连接地址
  ├── 执行测试
  └── 停止并清理容器

Testcontainers 不等于 Docker Compose,也不等于完整的生产环境模拟:

  • 它通常只启动测试所需要的依赖;
  • 容器端口可以映射到宿主机的动态端口;
  • 测试代码可以在启动后读取真实连接信息;
  • 它不会自动替你设计数据库清理、消息消费确认或业务断言;
  • 容器能启动不代表服务已经可以接受业务请求。

因此,应把 Testcontainers 看成一种测试环境编排工具,而不是测试断言框架。

二、从 Docker 镜像到业务断言的完整数据流

一个典型测试的状态转换如下:

sequenceDiagram
    participant P as pytest
    participant D as Docker/Testcontainers
    participant DB as PostgreSQL
    participant MQ as RabbitMQ/Kafka
    participant S as 被测服务

    P->>D: 创建数据库容器
    D->>DB: 启动 PostgreSQL
    P->>DB: readiness probe
    DB-->>P: 可连接、可执行查询

    P->>D: 创建消息队列容器
    D->>MQ: 启动 RabbitMQ/Kafka
    P->>MQ: readiness probe
    MQ-->>P: 可建立客户端连接

    P->>S: 注入数据库和消息队列连接配置
    P->>S: 发送业务请求
    S->>DB: 开启事务并写入业务数据
    S->>MQ: 发布领域事件
    S-->>P: 返回响应

    P->>DB: 查询最终数据
    P->>MQ: 消费或检查消息
    P-->>P: 断言响应、数据和消息副作用

    P->>D: 关闭服务、数据库和消息队列

这条路径中有三个不同的“成功”:

  1. 容器成功启动:Docker 创建并运行了容器。
  2. 服务成功就绪:数据库可以接受连接,消息代理可以完成协议握手。
  3. 业务操作成功:事务提交,消息发布,消费者完成处理。

把第一个状态误认为第三个状态,是最常见的集成测试错误。

三、环境准备:固定依赖,而不是固定宿主机

1. 安装测试依赖

下面的示例使用:

  • Python 3.14;
  • pytest
  • testcontainers
  • SQLAlchemy;
  • PostgreSQL 驱动;
  • pika 连接 RabbitMQ;
  • kafka-python-ng 连接 Kafka。

可以将测试依赖写入项目的测试依赖组:

[project]
requires-python = ">=3.14"

[dependency-groups]
test = [
    "pytest",
    "testcontainers[postgres,rabbitmq,kafka]",
    "sqlalchemy",
    "psycopg[binary]",
    "pika",
    "kafka-python-ng",
]

Testcontainers Python 当前项目将 postgresrabbitmqkafka 作为可选模块;RabbitMQ 模块使用 pika,Kafka 测试依赖中包含 Kafka Python 客户端。具体客户端版本应由项目锁文件管理,而不应在测试运行时随意解析最新版本。(github.com)

安装和运行:

python3.14 -m venv .venv
. .venv/bin/activate

python -m pip install -U pip
python -m pip install \
  pytest \
  "testcontainers[postgres,rabbitmq,kafka]" \
  sqlalchemy \
  "psycopg[binary]" \
  pika \
  kafka-python-ng

docker version
pytest -q

前置条件是当前用户能够访问 Docker daemon。常见失败不是 Python 代码错误,而是:

Cannot connect to the Docker daemon

此时应先验证:

docker info

docker info 能够返回 Server 部分,才说明 Testcontainers 具备创建容器的基础条件。

2. 固定镜像标签

示例中可以使用:

PostgresContainer("postgres:16")
RabbitMqContainer("rabbitmq:3.13-management")
KafkaContainer("confluentinc/cp-kafka:7.6.1")

镜像标签应当显式固定。使用 latest 会使测试结果随着镜像更新而变化,造成以下问题:

  • 今天通过,明天因为镜像升级失败;
  • 本地和 CI 拉到不同内容;
  • 数据库默认配置或消息代理协议行为发生变化;
  • 失败时无法恢复到原来的环境。

更严格的供应链控制可以使用镜像 digest,例如:

postgres@sha256:<digest>

但 digest 维护成本更高。版本标签和锁定的镜像清单之间需要做取舍:标签便于升级,digest 更利于复现。

四、数据库容器:连接、迁移和事务隔离

1. 一个最小 PostgreSQL 集成测试

Testcontainers 项目给出的基本模式是启动 PostgresContainer,调用 get_connection_url() 得到 SQLAlchemy 可用的连接 URL,再建立数据库连接执行查询。(github.com)

# tests/test_postgres.py
from sqlalchemy import create_engine, text
from testcontainers.postgres import PostgresContainer


def test_postgres_is_reachable():
    with PostgresContainer("postgres:16") as postgres:
        engine = create_engine(postgres.get_connection_url())

        with engine.connect() as connection:
            value = connection.execute(text("select 1")).scalar_one()

        engine.dispose()

    assert value == 1

每一步的含义是:

  1. PostgresContainer("postgres:16") 描述要运行的镜像;
  2. with 进入时启动容器并等待其连接信息可用;
  3. get_connection_url() 返回的是运行时 URL,主机端口可能每次不同;
  4. create_engine() 创建 SQLAlchemy 引擎;
  5. connection.execute() 通过真实 PostgreSQL 协议执行 SQL;
  6. engine.dispose() 关闭连接池;
  7. 离开 with 后容器停止并清理。

这里的断言只能说明 PostgreSQL 可以执行查询,不能说明迁移正确,也不能说明应用使用的连接池、事务和编码配置正确。

2. 在测试中执行迁移

应用启动时通常会通过 Alembic 或其他迁移工具创建表。集成测试应尽量执行与生产相同的迁移路径,而不是在测试代码中复制一份“简化建表 SQL”。

一种可测试的项目结构:

app/
  db.py
  models.py
  service.py
migrations/
tests/
  conftest.py
  test_order_repository.py

示例 fixture:

# tests/conftest.py
import os
import subprocess
from collections.abc import Iterator

import pytest
from sqlalchemy import create_engine
from sqlalchemy.engine import Engine
from testcontainers.postgres import PostgresContainer


@pytest.fixture(scope="session")
def postgres_url() -> Iterator[str]:
    with PostgresContainer("postgres:16") as postgres:
        yield postgres.get_connection_url()


@pytest.fixture(scope="session")
def db_engine(postgres_url: str) -> Iterator[Engine]:
    engine = create_engine(postgres_url, pool_pre_ping=True)

    environment = {
        **os.environ,
        "DATABASE_URL": postgres_url,
    }

    subprocess.run(
        ["alembic", "upgrade", "head"],
        check=True,
        env=environment,
    )

    yield engine
    engine.dispose()

这里使用 session 作用域,是因为启动 PostgreSQL 和执行迁移通常比创建一个测试事务更昂贵。pytest fixture 默认作用域是 function,也可以使用 classmodulepackagesession;fixture 会在对应作用域结束时销毁。(docs.pytest.org)

session 只适合共享基础设施,不代表共享业务数据。下面两个层次应分开:

session 级:PostgreSQL 容器、数据库连接配置、迁移结果
function 级:事务、测试数据、临时队列和断言上下文

3. 事务回滚不是万能隔离

很多团队使用如下方案:

@pytest.fixture
def db_session(db_engine):
    connection = db_engine.connect()
    transaction = connection.begin()

    yield connection

    transaction.rollback()
    connection.close()

如果被测代码使用的就是这个 connection,并且所有写操作都在该事务中执行,那么测试结束后可以回滚数据。

但是,这个方案有明确边界:

边界一:应用创建了自己的连接

测试 fixture 持有连接 A,而应用通过连接池创建连接 B:

测试事务:connection A
应用写入:connection B

A 的回滚不会撤销 B 已提交的事务。

边界二:应用提交了事务

如果应用调用:

connection.commit()

测试外部事务通常无法再回滚已经提交的数据。

边界三:异步消费者使用独立进程

消费者通过独立进程和独立连接读取数据库。测试主线程的事务未提交时,消费者可能根本看不到数据;测试主线程回滚时,消费者却可能已经处理了其他已提交副作用。

边界四:数据库外部副作用

数据库事务不能回滚:

  • 已发送到 RabbitMQ 的消息;
  • 已发送到 Kafka 的消息;
  • 调用第三方 HTTP 接口;
  • 文件系统中的文件;
  • 发送到另一个数据库的提交。

所以,事务回滚适合验证单个数据库连接内的同步读写,不应被误认为是整个分布式流程的隔离机制。

4. 三种数据库清理策略

策略 A:每个测试使用独立数据库

test_001 -> database_001
test_002 -> database_002

优点是隔离强,适合并行测试。缺点是创建数据库、运行迁移和管理连接的成本较高。

策略 B:共享数据库,测试后清空表

TRUNCATE TABLE orders, outbox_messages RESTART IDENTITY CASCADE;

优点是实现简单;缺点是如果测试失败、清理异常或新表没有加入清理列表,就会污染后续测试。

策略 C:共享数据库,每个测试使用唯一命名空间

例如所有表都有 tenant_id,每个测试生成随机租户:

from uuid import uuid4


def new_test_tenant() -> str:
    return f"test-{uuid4()}"

查询必须始终带上租户条件:

SELECT id, status
FROM orders
WHERE tenant_id = :tenant_id
ORDER BY id;

这种方案可以支持并行,但它要求业务代码和 SQL 查询都正确携带命名空间。只要有一条查询遗漏条件,隔离就会失效。

五、pytest Fixture:把生命周期显式化

1. yield fixture 的状态模型

一个 fixture 可以表示为:

未创建已创建测试使用已清理\text{未创建} \rightarrow \text{已创建} \rightarrow \text{测试使用} \rightarrow \text{已清理}

代码中的 yield 是生命周期分界点:

@pytest.fixture
def resource():
    handle = create_resource()
    try:
        yield handle
    finally:
        close_resource(handle)

如果创建成功,测试无论通过、失败还是抛出异常,finally 都会执行。pytest 的 yield fixture 会在测试完成后执行 yield 后面的清理代码;如果 fixture 在执行到 yield 之前就失败,pytest 不会执行其后的清理部分。(docs.pytest.org)

这也是为什么应把多个有副作用的动作拆成多个 fixture:

@pytest.fixture
def container():
    c = create_container()
    try:
        yield c
    finally:
        c.stop()


@pytest.fixture
def connection(container):
    conn = connect(container)
    try:
        yield conn
    finally:
        conn.close()

依赖关系是:

container -> connection -> test

清理顺序则相反:

test 完成 -> connection.close() -> container.stop()

如果把“启动容器、创建连接、声明队列、启动消费者”全部写在一个巨大的 fixture 中,中间任一步失败,都可能留下难以清理的部分状态。

2. 动态作用域

启动 Docker 容器可能需要较长时间。pytest 支持通过可调用对象动态决定 fixture 作用域,例如根据命令行参数在 sessionfunction 之间切换。(docs.pytest.org)

def container_scope(fixture_name, config):
    if config.getoption("--reuse-containers", default=False):
        return "session"
    return "function"


@pytest.fixture(scope=container_scope)
def postgres():
    with PostgresContainer("postgres:16") as container:
        yield container

对应命令行参数:

# tests/conftest.py
def pytest_addoption(parser):
    parser.addoption(
        "--reuse-containers",
        action="store_true",
        help="reuse a container during one pytest session",
    )

这个选项改变的是测试运行时的性能和隔离边界

  • 默认 function:测试间隔离更强,但启动次数更多;
  • session:运行更快,但所有测试共享服务状态,必须额外清理数据和消息。

“复用容器”不能等同于“复用测试数据”。

六、消息队列集成测试:发布成功不等于消费成功

消息队列测试比数据库测试更容易出现假阳性,因为消息处理通常是异步的。

设生产者发布消息 mm,消费者处理消息并更新数据库:

P(m)B(m)C(m)W(m)P(m) \rightarrow B(m) \rightarrow C(m) \rightarrow W(m)

其中:

  • PP:生产;
  • BB:消息代理暂存;
  • CC:消费者获取;
  • WW:业务写入数据库。

测试不能只验证:

producer.publish(message)
assert True

因为这只证明生产者函数没有立即抛异常。它没有证明:

  • 代理接受了消息;
  • 消息进入了正确的队列或 topic;
  • 消费者能够反序列化;
  • 消费者处理完成;
  • 数据库更新成功;
  • 消费确认发生在正确位置。

1. RabbitMQ:队列隔离和确认

RabbitMQ 使用队列保存消息。测试应为每个测试声明唯一队列,避免多个测试消费者读取同一批消息。

from uuid import uuid4

import pika
from testcontainers.rabbitmq import RabbitMqContainer


def test_rabbitmq_publish_and_consume():
    queue_name = f"orders.test.{uuid4().hex}"

    with RabbitMqContainer("rabbitmq:3.13-management") as rabbitmq:
        connection = pika.BlockingConnection(
            rabbitmq.get_connection_params()
        )
        channel = connection.channel()

        channel.queue_declare(queue=queue_name, durable=False)
        channel.basic_publish(
            exchange="",
            routing_key=queue_name,
            body=b'{"order_id": "o-1"}',
        )

        method, properties, body = channel.basic_get(
            queue=queue_name,
            auto_ack=False,
        )

        assert method is not None
        assert body == b'{"order_id": "o-1"}'

        channel.basic_ack(delivery_tag=method.delivery_tag)
        channel.queue_delete(queue=queue_name)
        connection.close()

关键点:

  • queue_name 使用随机后缀,避免测试间共享消息;
  • auto_ack=False 时,测试必须显式确认消息;
  • 断言消息后再 basic_ack,否则断言失败时消息可能已经被确认;
  • 最后删除测试队列,避免共享 RabbitMQ 容器时积累临时资源。

如果测试的是业务消费者,则更真实的流程应是:

发布消息
  ↓
等待消费者处理
  ↓
查询数据库状态
  ↓
断言状态变化

不要用固定的:

time.sleep(1)

因为它无法表达“什么时候算完成”。应使用带超时的轮询:

import time


def wait_until(predicate, timeout=10.0, interval=0.1):
    deadline = time.monotonic() + timeout

    while time.monotonic() < deadline:
        if predicate():
            return
        time.sleep(interval)

    raise AssertionError(f"condition not satisfied within {timeout}s")

使用方式:

wait_until(
    lambda: repository.get_status("o-1") == "processed",
    timeout=15,
)

它验证的是:

t[0,T],  status(t)="processed"\exists t \in [0, T],\; \text{status}(t) = \text{"processed"}

而不是假设固定时间后状态必然完成。超时 TT 必须有限,否则消费者挂死时测试会无限等待。

2. Kafka:topic、group 和 offset 隔离

Kafka 测试需要同时隔离:

  1. topic 名称
  2. consumer group ID
  3. offset 起始策略

如果多个测试使用同一个 group,Kafka 会把分区分配给不同消费者,导致某个测试收不到本应属于自己的消息。即使 topic 相同,只要 group 不同,消费者看到的消费进度也不同。

示例:

from uuid import uuid4

from kafka import KafkaConsumer, KafkaProducer
from testcontainers.kafka import KafkaContainer


def test_kafka_message_round_trip():
    topic = f"orders.test.{uuid4().hex}"
    group_id = f"group.test.{uuid4().hex}"

    with KafkaContainer("confluentinc/cp-kafka:7.6.1") as kafka:
        bootstrap = kafka.get_bootstrap_server()

        producer = KafkaProducer(
            bootstrap_servers=bootstrap,
            value_serializer=lambda value: value.encode(),
        )
        producer.send(topic, "order-created").get(timeout=10)
        producer.flush()
        producer.close()

        consumer = KafkaConsumer(
            topic,
            bootstrap_servers=bootstrap,
            group_id=group_id,
            auto_offset_reset="earliest",
            enable_auto_commit=False,
            value_deserializer=lambda value: value.decode(),
            consumer_timeout_ms=10_000,
        )

        records = list(consumer)
        consumer.close()

    assert [record.value for record in records] == ["order-created"]

这里设置 auto_offset_reset="earliest",是因为测试消费者可能在消息生产之后才启动;如果使用默认策略,消费者可能从最新位置开始,从而错过已经写入的测试消息。

consumer_timeout_ms 是有限等待的边界。没有超时,list(consumer) 可能持续等待;超时时间过短,则会把正常的 broker 初始化延迟误判为失败。

Kafka 的“消息已被消费”还不是业务处理完成。真实消费者通常在以下时刻提交 offset:

读取消息
  ↓
执行业务逻辑
  ↓
数据库事务提交
  ↓
提交 Kafka offset

如果先提交 offset,再提交数据库事务:

提交 offset
  ↓
数据库提交失败

消费者重启后不会再次读取该消息,业务数据就丢失。

如果先提交数据库,再提交 offset:

数据库提交成功
  ↓
offset 提交失败

消费者重启后可能重复处理,因此业务必须具备幂等性。

这形成了一个重要事实:

至少一次投递业务处理必须可重复\text{至少一次投递} \Rightarrow \text{业务处理必须可重复}

集成测试应专门覆盖:

  • 同一消息投递两次;
  • 数据库更新成功但 offset 提交失败;
  • 消费者处理失败后消息重新可见;
  • 消息格式错误时是否进入死信队列或被记录。

七、数据库与消息队列的一致性:测试 Outbox,而不是假设原子提交

应用经常需要同时完成:

写入 orders 表
发布 order-created 消息

数据库事务和消息代理通常属于两个独立资源,不能因为代码写在同一个函数里,就自动形成一个原子事务:

DB commit≢MQ publish\text{DB commit} \not\equiv \text{MQ publish}

存在两条失败路径。

路径一:先写数据库,再发消息

DB commit 成功
MQ publish 失败

结果是订单已存在,但下游没有收到事件。

路径二:先发消息,再写数据库

MQ publish 成功
DB commit 失败

结果是下游收到了一条描述不存在订单的消息。

Outbox 模式把消息写入同一个数据库事务:

BEGIN
  INSERT INTO orders ...
  INSERT INTO outbox_messages ...
COMMIT

之后由独立 relay 将 outbox_messages 发布到消息代理:

flowchart LR
    API[业务请求] --> TX[数据库事务]
    TX --> O[orders]
    TX --> B[outbox_messages]
    B --> R[Outbox Relay]
    R --> MQ[RabbitMQ/Kafka]
    MQ --> C[消费者]
    C --> DB2[下游数据库]

测试应分成两类:

  1. 事务测试:订单和 outbox 记录要么同时存在,要么同时不存在;
  2. relay 测试:outbox 记录最终会发布,发布失败会重试,成功后状态会更新。

例如验证事务原子性:

def test_order_and_outbox_are_committed_together(db_session):
    order_id = create_order_and_outbox(
        db_session,
        order_id="o-1",
        event_type="order-created",
    )

    order = db_session.execute(
        text("SELECT id FROM orders WHERE id = :id"),
        {"id": order_id},
    ).one()

    outbox = db_session.execute(
        text("""
            SELECT aggregate_id, event_type
            FROM outbox_messages
            WHERE aggregate_id = :id
        """),
        {"id": order_id},
    ).one()

    assert order.id == "o-1"
    assert outbox == ("o-1", "order-created")

真正的 relay 集成测试则需要同时启动数据库和消息队列,并验证消息最终进入指定队列或 topic。测试不能只检查 outbox 表,因为那只能证明数据库写入成功。

八、准备就绪:等待条件必须对应真实协议

容器状态有多个层次:

created   -> Docker 已创建
running   -> 主进程正在运行
ready     -> 服务接受连接并能完成最小协议交互
usable    -> 服务满足本测试的业务前置条件

例如:

  • PostgreSQL 的 running 不等于可以执行 SQL;
  • RabbitMQ 的进程启动不等于 AMQP 连接已可建立;
  • Kafka 的 broker 进程启动不等于 topic 元数据已经可用。

Testcontainers 的模块会提供相应的启动等待逻辑,但版本升级后等待策略可能变化。项目变更记录中曾包含 Kafka、RabbitMQ 等模块从日志等待迁移到更明确的等待策略,以及 RabbitMQ readiness probe 的修复。(github.com)

因此,测试代码仍应对业务级就绪负责。例如:

def wait_for_database(engine):
    wait_until(
        lambda: _can_query(engine),
        timeout=15,
    )


def _can_query(engine):
    try:
        with engine.connect() as connection:
            return connection.execute(text("select 1")).scalar_one() == 1
    except Exception:
        return False

这个探针比“容器日志出现某个字符串”更接近测试真正需要的条件:应用能否通过实际客户端执行查询。

日志等待也有风险:

  • 日志文本可能随镜像版本改变;
  • 日志出现不代表端口已可用;
  • 多节点服务可能在部分组件就绪时就输出类似文本;
  • 日志格式不是服务协议的一部分。

九、隔离的四个维度

“每次测试启动一个容器”只是隔离的一种实现。稳定隔离至少包含四个维度。

1. 网络隔离

如果两个测试共享一个 Docker network,它们可以通过相同服务名互相访问;如果测试需要模拟服务间通信,应显式创建网络并将相关容器加入同一网络。

但宿主机测试进程访问容器时,通常使用动态映射端口和 localhost;容器访问容器时,则使用 Docker network 中的服务名。两种地址不能混用:

宿主机 Python 进程 -> localhost:动态端口
容器内应用       -> postgres:5432

这也是 CI 中常见的“本地可用、容器内失败”原因。

2. 数据隔离

数据库隔离包括:

  • 独立容器;
  • 独立数据库;
  • 独立 schema;
  • 独立事务;
  • 唯一测试数据;
  • 测试后清理。

消息队列隔离包括:

  • 独立 queue;
  • 独立 topic;
  • 独立 consumer group;
  • 清空待处理消息;
  • 关闭消费者并等待其退出。

只隔离数据库而不隔离消息队列,仍可能发生跨测试污染:

测试 A 发布消息
测试 A 失败,未消费
测试 B 启动消费者并读到测试 A 的消息

3. 进程隔离

消费者可能是后台线程、子进程或独立容器。测试结束时,必须明确:

停止消费 -> 等待线程/进程退出 -> 关闭客户端 -> 删除队列或 topic

如果只关闭主测试函数,不停止后台消费者,后续测试可能继续修改数据库。

4. 时间隔离

异步系统的结果不是立即可见的。测试需要区分:

  • 同步断言:函数返回后即可检查;
  • 异步断言:必须等待某个状态最终出现;
  • 负向断言:在一段时间内不应出现某个副作用。

负向断言尤其困难:

assert no_message_arrives()

它无法在有限时间内证明“永远不会有消息”。可操作的定义是:

t[0,T],  ¬message_arrived(t)\forall t \in [0,T],\; \neg \text{message\_arrived}(t)

也就是在明确的观测窗口 TT 内没有消息到达。这个结论只对该窗口成立。

十、并行执行:隔离条件必须满足可组合性

假设测试 T1T_1T2T_2 并行运行。若它们共享状态 SS,则要求:

Δ(T1,S)Δ(T2,S)=\Delta(T_1, S) \cap \Delta(T_2, S) = \varnothing

其中 Δ(T,S)\Delta(T, S) 表示测试对状态 SS 的读写范围。

例如:

T1 写 orders.id = "o-1"
T2 写 orders.id = "o-1"

两者的写集合相交,并行结果依赖调度顺序,测试就不具备可组合性。

修复方式可以是:

from uuid import uuid4

order_id = f"order-{uuid4().hex}"
queue_name = f"queue.test.{uuid4().hex}"
group_id = f"group.test.{uuid4().hex}"

但随机 ID 只能避免名称碰撞,不能解决共享全局状态。例如:

  • 共享一个固定的 Redis key;
  • 共享一个固定 Kafka group;
  • 共享一个固定 RabbitMQ queue;
  • 依赖全局单例缓存;
  • 依赖同一个端口;
  • 测试修改了进程级环境变量。

pytest 支持通过标记选择测试,也可以注册自定义标记并启用严格标记检查,避免拼写错误导致测试筛选失效。(docs.pytest.org)

[tool.pytest.ini_options]
addopts = ["--strict-markers"]
markers = [
    "integration: requires Docker and external services",
    "serial: must not run in parallel",
]

运行:

pytest -m integration
pytest -m "integration and not serial"

serial 不应成为解决所有隔离问题的默认手段。它只是承认某些测试仍然共享不可分离的状态,代价是牺牲并行能力。更好的方向是先明确共享状态,再决定是否拆分容器、命名空间或测试数据。

十一、错误处理和诊断:先判断失败层次

1. 容器无法启动

表现:

DockerException
Image pull failed
port is already allocated
Cannot connect to the Docker daemon

诊断顺序:

docker info
docker ps -a
docker images
docker events

然后检查:

  • Docker daemon 是否运行;
  • CI runner 是否有 Docker socket 权限;
  • 镜像仓库是否可访问;
  • CPU 架构是否匹配;
  • 端口是否被固定绑定;
  • 是否配置了代理或私有 registry。

Testcontainers Python 使用 Ryuk 等机制帮助回收测试容器,也提供相关环境变量配置;禁用清理容器的机制会增加残留容器和磁盘占满的风险。(github.com)

2. 容器运行但服务未就绪

表现:

connection refused
timeout
leader not available
channel closed

此时查看容器日志:

docker ps
docker logs <container-id>

不要立即增加一个很大的 sleep。应确认:

  • 等待探针是否测试了正确端口;
  • 客户端连接的是映射端口还是容器端口;
  • Kafka 是否已经完成 broker 元数据初始化;
  • RabbitMQ 用户、vhost 和权限是否正确;
  • PostgreSQL 是否已经完成初始化脚本。

3. 业务断言失败

如果连接成功但业务断言失败,应保存:

  • 测试使用的镜像标签;
  • 容器日志;
  • 数据库迁移版本;
  • 发送消息的 topic、queue、group;
  • 消费者 offset 或 RabbitMQ ack 状态;
  • 相关 SQL 查询结果。

例如,失败时打印队列深度或数据库中的 outbox 状态,比只打印“expected processed, got pending”更有诊断价值。

十二、pytest 与 unittest 的边界

pytest fixture 适合表达资源依赖和生命周期:

@pytest.fixture
def service(db_engine, rabbitmq):
    ...

测试函数通过参数请求 fixture,pytest 根据依赖关系创建资源。pytest 文档将 fixture 描述为可组合、可按作用域共享的测试资源机制,并特别强调 teardown 和安全清理。(docs.pytest.org)

unittest 则提供标准库中的 TestCasesetUptearDownsetUpClasstearDownClass 等 xUnit 风格生命周期。Python 3.14 的官方文档仍将其作为标准测试框架提供。

在 pytest 中可以运行 unittest.TestCase,但两者有一个重要限制:unittest.TestCase 的测试方法不能直接接收 pytest fixture 参数。pytest 官方文档也指出,fixture 可以与 xUnit 风格 setup 混用,但 TestCase 方法不能通过参数注入 fixture。(docs.pytest.org)

因此,容器集成测试通常更适合写成 pytest 风格:

def test_order_flow(db_engine, rabbitmq):
    ...

如果已有 unittest 代码,可以将容器资源放入 setUpClass

import unittest

from testcontainers.postgres import PostgresContainer


class TestRepository(unittest.TestCase):
    @classmethod
    def setUpClass(cls):
        cls.postgres = PostgresContainer("postgres:16")
        cls.postgres.start()

    @classmethod
    def tearDownClass(cls):
        cls.postgres.stop()

    def test_connection(self):
        self.assertIsNotNone(self.postgres.get_connection_url())

这种写法能够工作,但资源依赖会隐藏在类属性中,无法像 pytest fixture 那样自然表达“数据库被哪些测试使用”“哪个资源先创建”“哪个资源后销毁”。大型集成测试套件中,pytest fixture 往往更容易组合和诊断。

十三、一个完整的订单事件测试结构

可以将一个真实流程拆成如下组件:

tests/
  conftest.py
  test_order_flow.py
app/
  api.py
  repository.py
  outbox.py
  consumer.py

测试步骤:

def test_create_order_emits_event_and_is_consumed(
    api_client,
    db_engine,
    message_consumer,
):
    response = api_client.post(
        "/orders",
        json={"order_id": "o-1", "amount": 100},
    )

    assert response.status_code == 201

    wait_until(
        lambda: read_order_status(db_engine, "o-1") == "processed",
        timeout=15,
    )

    event = message_consumer.read_one(timeout=5)
    assert event["type"] == "order-created"
    assert event["order_id"] == "o-1"

这个测试包含三种断言:

  1. HTTP 层:请求被接受;
  2. 数据库层:最终状态完成;
  3. 消息层:事件具有正确类型和业务主键。

如果只断言 HTTP 返回 201,测试可能在消息发布失败时仍然通过。反过来,如果只查询数据库,也无法证明事件已经发布给下游。

测试的输入、状态和副作用应尽量使用唯一标识:

from uuid import uuid4

order_id = f"o-{uuid4().hex}"

这样做的原因不是“随机更好”,而是使不同测试的状态集合尽量不相交:

{order_idT1}{order_idT2}=\{order\_id_{T_1}\} \cap \{order\_id_{T_2}\} = \varnothing

十四、CI 中的运行模型

CI 运行集成测试时,至少要明确以下结构:

Python 环境
  ├── 安装锁定的测试依赖
  ├── 拉取固定版本镜像
  ├── 启动 Testcontainers
  ├── 执行迁移和集成测试
  ├── 收集 junit、覆盖率和容器日志
  └── 清理容器与临时数据

常见命令:

python -m pytest -m integration \
  --junitxml=test-results/integration.xml \
  -vv

如果测试失败,CI 应保留:

docker ps -a
docker logs <database-container>
docker logs <message-broker-container>

但日志收集必须发生在容器清理之前。否则,失败原因可能随容器删除而消失。

依赖缓存和镜像缓存也要区分:

  • Python 包缓存减少安装时间;
  • Docker layer 缓存减少镜像拉取和构建时间;
  • Testcontainers 复用容器减少启动时间,但会削弱隔离;
  • 共享 CI runner 的 Docker daemon 可能包含其他作业的残留容器。

在多作业并行时,应避免固定宿主机端口。让 Testcontainers 使用动态端口映射,再通过容器对象提供的连接地址建立客户端连接,能够降低作业间端口冲突。

十五、常见误解和反例

误解一:容器启动成功就说明依赖可用

反例:

container.start()
client = Client("localhost", port)

start() 返回只说明 Testcontainers 完成了它的启动流程;业务客户端仍可能遇到服务初始化延迟。必须通过真实客户端完成最小连接或查询。

误解二:测试结束后停止容器,就不会有污染

反例:

共享 RabbitMQ 容器
测试 A 发布消息后失败
测试 A 没有删除 queue
测试 B 复用容器并消费到旧消息

容器生命周期和业务状态生命周期不是同一个概念。容器停止可以删除整个环境,但在共享容器期间必须自行清理队列、topic、数据库记录和消费者状态。

误解三:数据库回滚可以覆盖消息队列

反例:

BEGIN
INSERT order
COMMIT
publish message  # 失败
ROLLBACK         # 已无可回滚的数据库事务

回滚只能作用于仍未提交的那个数据库事务,不能撤销已经提交的数据库操作,也不能撤销消息代理中的发布。

误解四:固定等待时间比轮询简单

反例:

time.sleep(2)
assert status == "processed"

在本地机器上可能通过,在负载较高的 CI runner 上可能失败;将等待改成 10 秒又会拖慢本来只需 100 毫秒的测试。有限超时轮询同时表达了“最终一致性”和“不能无限等待”两个条件。

误解五:测试只需要覆盖成功路径

消息和数据库集成中,更有价值的失败路径包括:

数据库连接在提交前断开
消息发布超时
消费者反序列化失败
消费者处理成功但确认失败
消费者重复收到同一消息
迁移版本不匹配
容器启动后立即重启

这些路径决定了系统是丢数据、重复数据,还是能够安全重试。

十六、稳定隔离的判断标准

一个集成测试环境是否稳定,可以用以下条件判断:

环境可重建

测试不依赖开发者手工创建的数据库、固定本地端口或预先存在的队列。

状态可观测

失败时能够知道:

  • 容器是否启动;
  • 服务是否就绪;
  • 数据库迁移到哪个版本;
  • 消息发布到了哪里;
  • 消费者是否读取;
  • ack 或 offset 是否提交;
  • 数据库最终状态是什么。

副作用可清理

每一个成功的状态变更都有对应清理路径:

启动容器       -> 停止容器
创建连接       -> 关闭连接
创建队列       -> 删除队列
写入测试数据   -> 回滚、删除或销毁数据库
启动消费者     -> 停止并等待退出
创建临时文件   -> 删除临时目录

pytest 的 tmp_path 为每次测试调用提供唯一临时目录,适合保存迁移输出、消费者日志或导出的中间文件;该 fixture 的目录按测试调用隔离。(docs.pytest.org)

并行可组合

两个测试同时运行时,不会因为共享订单 ID、queue、topic、group、端口或全局环境变量而互相改变结果。

等待有边界

所有异步等待都满足:

0<T<0 < T < \infty

既不能无限等待,也不能使用没有业务依据的固定睡眠时间代替状态条件。

版本可追踪

测试结果能够关联到:

  • Python 3.14;
  • 锁定的 Python 依赖;
  • 固定的 Docker 镜像标签或 digest;
  • Docker daemon 版本;
  • CI runner 架构。

Testcontainers Python 的发布记录和项目配置会随版本变化;例如当前项目元数据与变更记录中已经包含 Python 3.14 分类和 4.15.0 版本信息。使用具体 API 前,应以项目当前版本的官方文档和锁文件为准,而不要把不同版本的示例混用。(github.com)

最终,Testcontainers 的价值不在于“让测试使用 Docker”,而在于把外部依赖的创建、就绪、连接、数据流、失败路径和销毁过程纳入测试控制。数据库测试要处理事务和迁移,消息队列测试要处理投递、确认、offset 和重复消费,pytest fixture 要处理资源依赖和反向清理;只有这些机制同时成立,集成测试才不仅是“偶尔能跑通”,而是能够在本地、CI 和并行执行中持续提供可信反馈。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。