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

Django Celery测试报错:connection already closed及任务执行失败

问题分析与解决方案

问题1:django.db.utils.InterfaceError: connection already closed

原因:Celery worker是独立进程,与测试进程的数据库连接不共享。默认的@pytest.mark.django_db会在测试结束后关闭数据库连接,但worker进程可能仍在复用这个已关闭的连接,触发报错。

问题2:任务状态为PENDING、数据库无数据

核心原因有三点:

  1. 字典键语法错误:测试代码中message_data的键未加引号(text: 'hello bro'),Python会将text视为变量而非字符串键,任务执行时触发NameError;但因任务标记了ignore_result=True,无法获取错误信息,状态显示为PENDING。
  2. 事务隔离:@pytest.mark.django_db(transaction=True)会让测试在事务中运行,Celery worker作为独立进程无法看到测试事务内的未提交数据,同时worker的写入操作也处于独立事务,测试进程无法感知。
  3. ignore_result=True的影响:该参数会阻止Celery保存任务结果,导致AsyncResult无法获取正确的任务状态。

具体修复步骤

1. 修复测试代码的语法错误

将message_data中的键改为字符串形式:

message_data = [
    {
        'text': 'hello bro',  # 补全字符串引号
        'client': 'Nick'     # 补全字符串引号
    },
]

2. 调整任务的结果保存设置

如果需要在测试中获取任务状态,暂时移除ignore_result=True(生产环境可根据需求恢复):

@shared_task  # 去掉ignore_result=True
def create_entries(data: list):
    batch_size = 100
    obj_iterator = (Message(**obj) for obj in data)
    while True:
        batch = list(islice(obj_iterator, batch_size))
        if not batch:
            break
        Message.objects.bulk_create(batch, batch_size)

3. 解决数据库连接与事务隔离问题

使用@pytest.mark.django_db(transaction=False)替代默认标记,避免事务隔离导致的进程间数据不可见问题;同时在任务中添加数据库连接清理逻辑,防止连接失效:

任务代码修改:

from .models import Message
from itertools import islice
from celery import shared_task
from django.db import close_old_connections

@shared_task
def create_entries(data: list):
    # 任务开始前关闭旧连接,避免连接已关闭报错
    close_old_connections()
    batch_size = 100
    obj_iterator = (Message(**obj) for obj in data)
    while True:
        batch = list(islice(obj_iterator, batch_size))
        if not batch:
            break
        Message.objects.bulk_create(batch, batch_size)
    # 任务结束后关闭连接
    close_old_connections()

测试代码修改:

import tasks
import pytest
from .models import Message
from celery.result import AsyncResult

@pytest.mark.django_db(transaction=False)  # 使用transaction=False
@pytest.mark.celery
def test_create_entries(celery_worker):
    message_data = [
        {
            'text': 'hello bro',
            'client': 'Nick'
        },
    ]
    assert Message.objects.count() == 0
    task = tasks.create_entries.delay(message_data)
    # 等待任务执行完成
    task.wait()
    result = AsyncResult(task.task_id)
    assert result.status == 'SUCCESS'
    assert Message.objects.count() == 1

4. 确保Celery Worker使用测试配置

在pytest.ini中指定测试设置,让celery_worker fixture加载test_settings.py:

[pytest]
DJANGO_SETTINGS_MODULE = your_project.test_settings

内容的提问来源于stack exchange,提问作者Sergey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 19:31:05