Celery任务串行执行求助:Django+Windows环境并行失效
问题背景
我有一个IO密集型的Django应用,用Celery搭配gevent运行任务,通过UI进度条管理任务进度。环境配置如下:
- Django版本:5.0.2
- Celery版本:5.3.6
- Redis版本:Windows版5.0.14.1
- 服务器:Windows Server 2016(无法更换,数据存储在Access数据库),IIS默认应用池部署,4核4GB内存
关键配置
web.config
<?xml version="1.0" encoding="utf-8"?> <configuration> <system.webServer> <handlers> <add name="Python FastCGI" path="*" verb="*" modules="FastCgiModule" scriptProcessor="C:\Python311\python.exe|C:\Python311\Lib\site-packages\wfastcgi.py" resourceType="Unspecified" requireAccess="Script" /> </handlers> <directoryBrowse enabled="true" /> </system.webServer> <appSettings> <add key="PYTHONPATH" value="C:\inetpub\Django-LIAL\WEBAPPLIAL" /> <add key="WSGI_HANDLER" value="WEBAPPLIAL.wsgi.application" /> <add key="DJANGO_SETTINGS_MODULE" value="WEBAPPLIAL.settings" /> </appSettings> </configuration>
Django WSGI配置(wsgi.py)
from gevent import monkey monkey.patch_all() import os from django.core.wsgi import get_wsgi_application os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'WEBAPPLIAL.settings') application = get_wsgi_application()
Celery配置(settings.py)
# Celery setting CELERY_BROKER_URL = 'redis://127.0.0.1:6379/0' CELERY_ACCEPT_CONTENT = ['json'] CELERY_TASK_SERIALIZER = 'json' CELERY_RESULT_BACKEND = 'django-db' CELERY_CACHE_BACKEND = 'django-cache' CELERY_TASK_ALWAYS_EAGER = False CELERY_TASK_TRACK_STARTED = True
Celery启动命令
celery -A WEBAPPLIAL worker -l info -P gevent
问题
Celery启动日志显示并发数为4,但通过.delay()调用的两个操作不同Django ORM模型的任务始终串行执行,无法并行处理。
解决方案
1. 给Celery Worker进程打gevent猴子补丁
当前仅在WSGI应用中打了gevent补丁,但Celery Worker是独立进程,必须单独补丁才能让gevent协程正确处理IO阻塞。
在Celery实例初始化文件(如celery.py)中添加补丁代码,确保在Celery初始化前执行:
from gevent import monkey monkey.patch_all() import os from celery import Celery os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'WEBAPPLIAL.settings') app = Celery('WEBAPPLIAL') app.config_from_object('django.conf:settings', namespace='CELERY') app.autodiscover_tasks()
2. 解决Access数据库的单连接阻塞问题
Access是单文件型数据库,默认驱动会锁定文件导致多请求阻塞,这是任务串行的核心原因之一:
- 设置短生命周期数据库连接:在
settings.py的数据库配置中添加CONN_MAX_AGE=0,禁用连接复用,每次任务都创建新连接:DATABASES = { 'default': { # 替换为你的Access数据库引擎(如django_accessdb或pyodbc) 'ENGINE': 'django_accessdb.backends.access', 'NAME': r'C:\path\to\your\database.accdb', 'CONN_MAX_AGE': 0, # 若用pyodbc,添加共享访问参数 'OPTIONS': { 'driver': 'Microsoft Access Driver (*.mdb, *.accdb)', 'extra_params': 'Mode=ReadWrite;Share Deny None', } } } - 任务内显式管理连接:如果任务涉及长时间IO操作,可在任务完成后手动关闭数据库连接:
from django.db import connection @shared_task def io_bound_task(): try: # 执行数据库操作 pass finally: connection.close()
3. 显式设置gevent并发数
IO密集型任务适合更高的并发数,默认的4核对应并发数不足以充分利用资源,启动命令改为:
celery -A WEBAPPLIAL worker -l info -P gevent --concurrency 15
可根据服务器内存调整并发数(建议10-20)。
4. 验证任务并行状态
在任务中添加日志,确认是否为不同协程处理:
import time import logging from celery import shared_task import gevent logger = logging.getLogger(__name__) @shared_task def task_one(): logger.info(f"Task 1启动 | 时间戳: {time.time()} | 协程ID: {id(gevent.getcurrent())}") # 模拟IO操作 time.sleep(5) logger.info(f"Task 1完成 | 时间戳: {time.time()}") @shared_task def task_two(): logger.info(f"Task 2启动 | 时间戳: {time.time()} | 协程ID: {id(gevent.getcurrent())}") # 模拟IO操作 time.sleep(5) logger.info(f"Task 2完成 | 时间戳: {time.time()}")
调用task_one.delay()和task_two.delay()后,若日志显示两个任务启动时间接近、协程ID不同,则说明并行生效。
内容的提问来源于stack exchange,提问作者Mougnou
相关产品推荐
相关产品推荐

