【问题标题】:django + celery: disable prefetch for one worker, Is there a bug?django + celery:为一名工作人员禁用预取,有错误吗?
【发布时间】: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

【问题讨论】:

    标签: python django celery


    【解决方案1】:

    我确认该问题,文档的this 部分存在错误。 worker_prefetch_multiplier = 1 就像它说的那样,将 worker 的 prefetch 设置为 1,这意味着 worker 将在当前正在执行的任务之外再持有一个任务。

    要真正禁用预取,您还需要使用 task_acks_late = True 和预取设置,请参阅 this 文档部分

    【讨论】:

    • 非常感谢。这会奏效。希望医生澄清。可惜的是,task_acks_late = True 不能设置为 celery 命令行选项。所以我必须为所有工作人员更改此设置,这并不是我真正想要的。我现在正在做的解决方法是有一个 settings.py 文件,根据环境变量的存在将task_acks_late 设置为TrueFalse,我调用celery ... worker 一次并设置环境变量和一次未设置环境变量。这真的很尴尬。任何建议如何以更直观的方式做到这一点。
    • 不客气,我认为这个新问题值得单独提出一个问题
    • 是的,我也想过,但不确定。所以我认为。我先发表评论。
    • 我创建了一个后续问题:stackoverflow.com/questions/58328194/…
    猜你喜欢
    • 2021-09-08
    • 2021-08-15
    • 2021-08-18
    • 2013-02-18
    • 2013-02-16
    • 1970-01-01
    • 1970-01-01
    • 2020-01-28
    • 1970-01-01
    相关资源
    最近更新 更多