用 pytest 隔离测试 Celery 任务和内嵌 worker

用 pytest 隔离测试 Celery 任务和内嵌 worker

Celery 测试需要区分任务边界的单元测试、带 worker 的集成测试,以及接近部署环境的冒烟测试。本文完整整理官方 Testing with Celery 一章;核验时稳定文档为 Celery 5.6.3。所有示例仅作静态审核。

单元测试用 mock 隔离业务依赖;集成测试经测试应用和专用 broker 到内嵌 worker,再由结果后端返回;fixture 生命周期应保持一致。
未完纪原创技术示意图,根据本文引用的官方文档绘制;非原站截图。

先区分两套测试接口

官方把 Celery 测试分成两部分:单元测试与集成测试使用 celery.contrib.pytest;冒烟测试和生产相关测试使用 pytest-celery >= 1.0.0。安装 pytest-celery 也会带入 celery.contrib.pytest 的基础设施,但两套 API 不兼容。celery.contrib.pytest 是 mock-based;pytest-celery 使用 Docker 测试环境。本文下列 fixture 示例只适用于前一套接口。Docker 型测试应单独阅读 pytest-celery 文档,不能混用下面的示例。

任务与单元测试

测试任务行为时,推荐使用 mock。任务很像 Web 视图:主要定义以任务方式调用时需要的处理,例如序列化、消息头和重试;真正的业务逻辑最好放到其他位置。

eager 模式的限制:task_always_eager 只模拟 worker 行为,与真实 worker 存在差异,不能替代消息传输、序列化和 worker 测试。eager 任务默认不向结果后端保存结果;需要这一行为时查看 task_store_eager_result。

以下沿用原文的下单场景。模型只传主键,不序列化整个对象。绑定任务的 bind=True 会把任务实例作为第一个参数 self 传入,因此可以调用 retry 等任务方法。编辑修订:补齐 Decimal、OperationalError 和项目 app 的导入上下文,并将对象查询修正为明确的关键字参数。项目需已有对应 Product 模型与 Celery app。

from decimal import Decimal
from django.db import OperationalError
from .celery import app
from .models import Product

@app.task(bind=True)
def send_order(self, product_pk, quantity, price):
    price = Decimal(price)  # JSON 金额使用字符串
    product = Product.objects.get(pk=product_pk)
    try:
        product.order(quantity, price)
    except OperationalError as exc:
        raise self.retry(exc=exc)

重试范围与幂等性:self.retry 只包围 product.order。如果 Product.objects.get 自身抛出 OperationalError,这个示例不会重试该查询。product.order 出错后也可能被再次调用;真实任务应先确认事务边界和重复执行是否安全,必要时使用幂等键。本示例没有证明这些条件,也没有扩大重试范围。

原文单元测试直接创建 Product,再替换 Product.order,因此仍会访问数据库,需要相应数据库 fixture。下面改为同时 mock ORM 查询,让这个版本不访问真实数据库。pytest 默认收集 Test 开头的类,原文小写类名也一并修正。金额以字符串构造,避免 Decimal(30.3) 把浮点误差带入样本。

from decimal import Decimal
from unittest.mock import Mock, patch
import pytest
from celery.exceptions import Retry
from django.db import OperationalError
from proj.tasks import send_order

class TestSendOrder:
    @patch("proj.tasks.Product.objects.get")
    def test_success(self, product_get):
        product = Mock()
        product_get.return_value = product
        send_order(1, 3, "30.3")
        product_get.assert_called_once_with(pk=1)
        product.order.assert_called_once_with(3, Decimal("30.3"))

    @patch("proj.tasks.Product.objects.get")
    @patch("proj.tasks.send_order.retry")
    def test_failure(self, send_order_retry, product_get):
        error = OperationalError("temporary database failure")
        product_get.return_value.order.side_effect = error
        send_order_retry.side_effect = Retry()
        with pytest.raises(Retry):
            send_order(1, 3, "30.6")
        send_order_retry.assert_called_once_with(exc=error)

将 proj 替换成实际项目包名。mock 应放在被测模块查找符号的位置;它不会自动隔离其他网络调用、信号处理器或全局状态。这里检查成功分支的金额转换,以及业务错误发生时任务是否请求重试。原文针对 Python 2 提到第三方 mock 包;本稿采用 Python 3 标准库 unittest.mock,不再提供过时环境安装步骤。

启用 pytest 插件

Celery 从 4.0 开始提供 pytest 插件,增加可供集成测试和单元测试使用的 fixture。插件默认未启用。官方给出三种启用方式:安装 celery[pytest];设置环境变量 PYTEST_PLUGINS=celery.contrib.pytest;或在根 conftest.py 中声明插件。按项目安装方式选择,避免重复注册。

pip install "celery[pytest]"
# 根目录 conftest.py
pytest_plugins = ("celery.contrib.pytest",)

celery mark:覆盖测试应用配置

celery mark 可以覆盖单个测试或整个测试类的配置。下例保留原文 Redis 后端短 URL 的示意;实际应指定专用测试资源,避免连到共享或生产 Redis。

import pytest

@pytest.mark.celery(result_backend="redis://")
def test_something():
    ...

@pytest.mark.celery(result_backend="redis://")
class TestSomething:
    def test_one(self):
        ...

    def test_two(self):
        ...

函数级 fixture

celery_app:当前测试使用的应用

这个 fixture 返回测试用 Celery app。任务在测试内动态注册后,需要重新加载 worker。下面的乘法例子还需要可用的 broker 和结果后端。

def test_create_task(celery_app, celery_worker):
    @celery_app.task
    def mul(x, y):
        return x * y

    celery_worker.reload()
    assert mul.delay(4, 4).get(timeout=10) == 16

celery_worker:内嵌一个真实 worker

fixture 在独立线程启动 worker,测试返回后关闭。默认最多等待 10 秒让未完成任务结束,超时会抛异常。可在 celery_worker_parameters 返回的字典中调整 shutdown_timeout;它与 AsyncResult.get 的结果读取超时控制不同阶段。

# conftest.py;URL 应改为专用测试资源
import pytest

@pytest.fixture(scope="session")
def celery_config():
    return {
        "broker_url": "amqp://",
        "result_backend": "redis://",
    }

# mytask 必须已注册在测试应用中
def test_add(celery_worker):
    mytask.delay()

@pytest.mark.celery(result_backend="rpc://")
def test_other(celery_worker):
    ...

worker 默认关闭 heartbeat,不会发送 worker-online、worker-offline 和 worker-heartbeat 事件。需要测试这些事件时,覆盖 worker 参数。以下把 heartbeat 与退出等待放在同一个 fixture 中。

@pytest.fixture(scope="session")
def celery_worker_parameters():
    return {
        "without_heartbeat": False,
        "shutdown_timeout": 20,
    }

编辑补充:同进程线程 worker 的小型测试也可采用 memory:// broker 和 cache+memory:// 后端,但它们不覆盖真实网络故障,不能跨进程共享消息或结果。不能将这一配置直接套到 prefork 环境。

会话级 fixture

celery_config 与 celery_parameters

重新定义 celery_config 可以配置测试应用,返回值同时用于 celery_app 和 celery_session_app。celery_parameters 则把参数直接传给 Celery 的构造函数,用于自定义任务基类、strict_typing 等初始化选项。

@pytest.fixture(scope="session")
def celery_config():
    return {
        "broker_url": "amqp://",
        "result_backend": "rpc://",
    }

@pytest.fixture(scope="session")
def celery_parameters():
    return {
        "task_cls": my.package.MyCustomTaskClass,
        "strict_typing": False,
    }

自定义任务类是项目占位符,必须按实际路径导入。一个 conftest.py 中只应保留一份最终采用的同名 fixture。这里将原文 rpc 简写写成完整的 rpc://,没有验证实际 broker 配置。

celery_worker_parameters

这个 fixture 控制测试 worker 的构造参数,直接交给 WorkController,对函数级与会话级 worker 都生效。例如只消费两个指定队列并排除默认队列:

@pytest.fixture(scope="session")
def celery_worker_parameters():
    return {
        "queues": ("high-prio", "low-prio"),
        "exclude_queues": ("celery",),
    }

编辑修订:原文的 (‘celery’) 是字符串;单元素元组必须加尾逗号。测试队列应使用隔离命名,避免消费者取走其他环境的消息。如还需 heartbeat 或 shutdown_timeout,应合并进同一个返回字典。

celery_enable_logging

覆盖这个 fixture 可启用内嵌 worker 的日志。

@pytest.fixture(scope="session")
def celery_enable_logging():
    return True

celery_includes

返回 worker 启动时需要额外导入的模块名列表,可用于任务模块和注册信号处理器的模块。

@pytest.fixture(scope="session")
def celery_includes():
    return [
        "proj.tests.tasks",
        "proj.tests.celery_signal_handlers",
    ]

celery_worker_pool

覆盖此 fixture 可选择内嵌 worker 的执行池,原文用 prefork 举例:

@pytest.fixture(scope="session")
def celery_worker_pool():
    return "prefork"

除非整个测试套件已在合适时机启用 monkeypatch,否则不能使用 gevent/eventlet 池。prefork 的平台支持、进程隔离、数据库连接和事务可见性也要独立核验。测试事务中尚未提交的数据,不应假定另一线程或进程都能读取。

celery_session_worker 与 celery_session_app

celery_session_worker 在整个会话中持续运行,无需每个测试启停;celery_session_app 供其他会话级 fixture 引用同一个 Celery 应用。可搭配前述 session 级 celery_config:

def test_add_task(celery_session_worker):
    # add 必须已注册到会话应用中
    assert add.delay(2, 2).get(timeout=10) == 4

本稿补了结果读取超时,防止配置错误导致无限等待。不要混用会话 worker 与短期函数 worker;共享生命周期会扩大状态串扰范围,需要在测试之间清理业务状态。

use_celery_app_trap:捕获隐式默认应用

在 conftest.py 中启用 app trap 后,代码若尝试访问默认应用或 current_app,就会抛异常。有意访问默认应用的测试可声明 depends_on_current_app fixture。

@pytest.fixture(scope="session")
def use_celery_app_trap():
    return True

@pytest.mark.usefixtures("depends_on_current_app")
def test_something():
    something()

something 是原文的操作占位符。app trap 针对应用选择错误,不会替你隔离文件、数据库和所有外部副作用;这些仍需要在测试设计中处理。

来源、署名与许可

来源:Celery 官方用户指南 Testing with Celery,稳定文档页标示 Celery 5.6.3。作者署名为 Celery 文档贡献者。文档内容及其中的教程示例按 CC BY-SA 4.0 授权;Celery 软件本身另按 BSD-3-Clause 许可。完整软件与文档许可文本、版权声明见随稿 LICENSE-Celery.txt。本文对示例作的结构、语法和隔离边界修订均已在相应位置说明。

本文中的命令和返回值是教程示例,并非本稿运行结果。未执行任务、测试、安装或连接 broker、后端、worker 与数据库。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容