【发布时间】:2020-02-05 23:00:37
【问题描述】:
我有一个带有 celery 的 Django 项目
由于内存限制,我只能运行两个工作进程。
我有“慢”和“快”的混合任务。 快速任务应尽快执行。在很短的时间内(0.1s - 3s)可以有很多快速的任务,所以理想情况下两个 CPU 都应该处理它们。
慢速任务可能会运行几分钟,但结果可能会延迟。
慢任务发生的频率较低,但可能会同时排队 2 或 3 个。
我的想法是拥有一个:
- 1 个 celery worker W1 并发 1,只处理快速任务
- 1 个具有并发 1 的 celery worker W2 可以处理快速和慢速任务。
默认情况下,celery 的任务预取乘数 (https://docs.celeryproject.org/en/latest/userguide/configuration.html#worker-prefetch-multiplier) 为 4,这意味着 4 个快速任务可以排在慢速任务之后,并且可能会延迟几分钟。因此,我想禁用工人 W2 的预取。该文档指出:
要禁用预取,请将 worker_prefetch_multiplier 设置为 1。 设置为 0 将允许工人继续消耗尽可能多的 消息随心所欲。
但是我观察到的是,prefetch_multiplier 为 1 时,一个任务会被预取,但仍会被慢速任务延迟。
这是一个文档错误吗?这是一个实现错误吗?还是我误解了文档? 有什么方法可以实现我想要的吗?
我为启动工作程序而执行的命令是:
celery -A miniclry worker --concurrency=1 -n w2 -Q=fast,slow --prefetch-multiplier 0
celery -A miniclry worker --concurrency=1 -n w1 -Q=fast
我的 celery 设置是默认设置,除了:
CELERY_BROKER_URL = "pyamqp://*****@localhost:5672/mini"
CELERY_TASK_ROUTES = {
'app1.tasks.task_fast': {"queue": "fast"},
'app1.tasks.task_slow': {"queue": "slow"},
}
我的 django 项目的 celery.py 文件是:
from __future__ import absolute_import
import os
from celery import Celery
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'miniclry.settings')
app = Celery("miniclry", backend="rpc", broker="pyamqp://")
app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks()
我的 django 项目的__init__.py 是
from .celery import app as celery_app
__all__ = ('celery_app',)
我的工人的代码
import time, logging
from celery import shared_task
from miniclry.celery import app as celery_app
logger = logging.getLogger(__name__)
@shared_task
def task_fast(delay=0.1):
logger.warning("fast in")
time.sleep(delay)
logger.warning("fast out")
@shared_task
def task_slow(delay=30):
logger.warning("slow in")
time.sleep(delay)
logger.warning("slow out")
如果我从我看到的管理 shell 中执行以下操作,那么只有在慢速任务完成后才会执行一个快速任务。
from app1.tasks import task_fast, task_slow
task_slow.delay()
for i in range(30):
task_fast.delay()
有人可以帮忙吗?
如果认为有帮助,我可以发布整个测试项目。只是建议交换此类项目的推荐 SO 方式
版本信息:
- celery==4.3.0
- Django==1.11.25
- Python 2.7.12
【问题讨论】: