动态添加周期性任务 celery

gga*_*par 8 python celery flask

是否可以向 celery动态添加周期性任务?

我正在使用 Flask,而不是 django,并且我正在构建一个应用程序,该应用程序应该允许用户通过 Web 界面定义循环任务。

我已经尝试使用 Celery 4.1 中的 Periodic Tasks,但是要添加新任务,我必须停止 celery 服务器,更改配置(即使通过 python 完成),然后重新启动它。也许有一种方法可以动态加载配置(无需重新启动)?

我考虑过有一个 crontab,每 5 分钟重新启动一次芹菜服务。但这似乎非常违反自然。在其他原因中,我想使用 celery 的原因是不使用 crontab。

有没有人对此有所了解?

ps。:我知道另一个类似的问题,但它是从 2012 年开始的。我希望从那时起事情就发生了变化,即在 v4.1 中引入了beat

ham*_*mSh 5

为此,您可以使用 redBeat。基于redbeat GitHub:

RedBeat 是一个 Celery Beat Scheduler,它将计划任务和运行时元数据存储在 Redis 中。

对于任务创建:

import tasks #celery defined task class 
from redbeat import RedBeatSchedulerEntry as Entry
entry = Entry(f'urlCheck_{key}', 'tasks.urlSpeed', repeat, args=['GET', url, timeout, key], app=tasks.app)
entry.save()
entry.key
Run Code Online (Sandbox Code Playgroud)

删除任务:

import tasks #celery defined task class 
from redbeat import RedBeatSchedulerEntry as Entry
entry = Entry.from_key(key, app=tasks.app) #key from previous step
entry.delete()
Run Code Online (Sandbox Code Playgroud)

有一个可以使用的示例:https ://github.com/hamedsh/redBeat_example


ana*_*nan 4

这适用于 Celery 4.0.1+ 和 Python 2.7 以及 Redis

from celery import Celery
import os, logging
logger = logging.getLogger(__name__)
current_module = __import__(__name__)

CELERY_CONFIG = {
    'CELERY_BROKER_URL': 
     'redis://{}/0'.format(os.environ.get('REDIS_URL', 'localhost:6379')),
  'CELERY_TASK_SERIALIZER': 'json',
}


celery = Celery(__name__, broker=CELERY_CONFIG['CELERY_BROKER_URL'])
celery.conf.update(CELERY_CONFIG)
Run Code Online (Sandbox Code Playgroud)

我通过以下方式定义工作:

job = {
    'task': 'my_function',               # Name of a predefined function
    'schedule': {'minute': 0, 'hour': 0} # crontab schedule
    'args': [2, 3],
    'kwargs': {}
}
Run Code Online (Sandbox Code Playgroud)

然后我定义一个这样的装饰器:

def add_to_module(f):
    setattr(current_module, 'tasks_{}__'.format(f.name), f)
    return f
Run Code Online (Sandbox Code Playgroud)

我的任务是

@add_to_module
def my_function(x, y, **kwargs):
    return x + y
Run Code Online (Sandbox Code Playgroud)

然后添加一个动态添加任务的函数

def add_task(job):
    logger.info("Adding periodic job: %s", job)
    if not isinstance(job, dict) and 'task' in jobs:
        logger.error("Job {} is ill-formed".format(job))
        return False
    celery.add_periodic_task(
        crontab(**job.get('schedule', {'minute': 0, 'hour': 0})),
        get_from_module(job['task']).s(
            enterprise_id,
            *job.get('args', []),
            **job.get('kwargs', {})
        ),
        name = job.get('name'),
        expires = job.get('expires')
    )
    return True


def get_from_module(f):
    return getattr(current_module, 'tasks_{}__'.format(f))
Run Code Online (Sandbox Code Playgroud)

之后,您可以将 add_task 函数链接到 URL,并让它们从当前模块中的函数创建任务