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
为此,您可以使用 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
这适用于 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,并让它们从当前模块中的函数创建任务
| 归档时间: |
|
| 查看次数: |
4651 次 |
| 最近记录: |