Celery pytest fixtures在HTTP请求测试中无法同步等待异步任务执行的问题咨询
嗨,我来帮你解决这个问题~ 你遇到的情况其实很常见:直接调用任务时你能拿到AsyncResult实例并调用get()等待,但视图里触发的任务没法直接拿到这个实例,所以测试跑断言的时候任务还在异步执行中。
不用task_always_eager也能搞定,这里有几个实用的方案:
方案一:利用测试worker的wait()方法(最推荐)
Celery的celery_session_worker fixture提供的worker实例有个wait()方法,它会阻塞直到worker完成所有队列中的任务。完全不用改业务代码,只需要在测试里加一行就行:
from django.core import mail from django.test import TestCase from django.urls import reverse import pytest @pytest.mark.usefixtures('celery_session_app') @pytest.mark.usefixtures('celery_session_worker') class SendMailTestCase(TestCase): def test_request_url(self, celery_session_worker): response = self.client.get(reverse('groups')) self.assertEqual(response.status_code, 200) # 等待worker完成所有待执行任务 celery_session_worker.wait() self.assertEqual(len(mail.outbox), 1)
这个方法的好处是完全不侵入业务逻辑,只在测试层处理,非常干净。
方案二:让视图返回任务ID,测试中主动等待
如果需要更精准地等待特定任务,可以修改视图在测试环境下返回任务ID,然后在测试里拿到ID后调用get()等待:
首先修改视图(用环境变量控制只在测试时返回):
import os from rest_framework.response import Response # ... 其他导入代码 class GroupView(ListAPIView): # ... 其他视图代码 def get(self, request, *args, **kwargs): task = task_send_mail.delay( email_address='mail@mail.org.br', message='hello' ) # 仅测试环境返回task_id if os.environ.get('TESTING') == 'True': return Response(status=200, data={'task_id': task.id}) return Response(status=200, data={})
然后修改测试代码:
from django.core import mail from django.test import TestCase from django.urls import reverse from polls.tasks import task_send_mail import pytest import os @pytest.mark.usefixtures('celery_session_app') @pytest.mark.usefixtures('celery_session_worker') class SendMailTestCase(TestCase): def setUp(self): # 设置测试环境标记 os.environ['TESTING'] = 'True' def tearDown(self): # 清理环境变量 del os.environ['TESTING'] def test_request_url(self): response = self.client.get(reverse('groups')) self.assertEqual(response.status_code, 200) # 获取任务ID并等待执行完成 task_id = response.data['task_id'] task = task_send_mail.AsyncResult(task_id) task.get() self.assertEqual(len(mail.outbox), 1)
这个方案适合需要单独验证某个任务的场景,但需要少量修改视图代码。
方案三:用Celery inspect API收集任务并等待
如果你不想改视图也不想依赖worker的wait()方法,可以用Celery的inspect API获取队列中所有待执行/正在执行的任务,然后逐个等待完成:
from django.core import mail from django.test import TestCase from django.urls import reverse from celery.result import AsyncResult import pytest @pytest.mark.usefixtures('celery_session_app') @pytest.mark.usefixtures('celery_session_worker') class SendMailTestCase(TestCase): def test_request_url(self, celery_session_app): response = self.client.get(reverse('groups')) self.assertEqual(response.status_code, 200) # 获取worker的inspect实例 inspect = celery_session_app.control.inspect() # 获取所有待执行和正在执行的任务 pending_tasks = inspect.pending() or {} active_tasks = inspect.active() or {} # 收集所有任务ID task_ids = [] for queue_tasks in pending_tasks.values(): task_ids.extend(task['id'] for task in queue_tasks) for queue_tasks in active_tasks.values(): task_ids.extend(task['id'] for task in queue_tasks) # 逐个等待任务完成 for task_id in task_ids: AsyncResult(task_id).get() self.assertEqual(len(mail.outbox), 1)
这个方案更灵活,但代码稍复杂,适合需要自定义等待逻辑的场景。
为什么原来的fixture没自动等待?
其实celery_session_app和celery_session_worker只是帮你启动了一个可用的Celery worker实例,让异步任务能被执行,但它们不会自动拦截或等待所有任务——毕竟在实际生产中,任务就是异步执行的,测试fixture只是模拟这个环境,不会主动改变任务的执行模式。task_always_eager是强制任务同步执行,相当于把异步任务变成了同步函数,而上面的方案是在测试里主动等待任务完成,更贴近真实的运行场景。
备注:内容来源于stack exchange,提问作者Charles Rockenbach

