一、背景实际工作中会有一些耗时的异步任务需要使用定时调度,比如发送邮件,拉取数据,执行定时脚本
通过celery 实现调度主要思想是 通过引入中间人redis,启动 worker 进行任务执行 ,celery-beat进行定时任务数据存储
二、celery动态添加定时任务的官方文档celery文档:https://docs.celeryproject.org/en/latest/userguide/periodic-tasks.html#beat-custom-schedulers
celery 自定义调度类说明:
自定义调度器类可以在命令行中指定(--scheduler参数)
django-celery-beat文档 : https://pypi.org/project/django-celery-beat/
关于django-celery-beat 插件的说明:
此扩展使您能够将定期任务计划存储在数据库中,可以从 django 管理界面管理周期性任务,您可以在其中创建、编辑和删除周期性任务以及它们应该运行的频率
三、celery简单实用3.1 基础环境配置1. 安装最新版本的django
pip3 install django #当前我安装的版本是 3.0.6
2. 创建项目
django-admin startproject typeideadjango-admin startapp blog
3.安装 celery
pip3 install django-celerypip3 install -u celery pip3 install "celery[librabbitmq,redis,auth,msgpack]" pip3 install django-celery-beat # 用于动态添加定时任务pip3 install django-celery-resultspip3 install redis
3.2 测试使用celery应用1. 创建blog目录、新建task.py
首先在django项目中创建一个blog文件夹,并且在blog文件夹下创建tasks.py模块, 如下:
tasks.py代码如下:
#!/usr/bin/env python# -*- coding: utf-8 -*- """#file: tasks.py#time: 2022/3/30 2:26 下午#author: julius"""from celery import celery # 使用redis做为brokerapp = celery('blog.tasks2',broker='redis://127.0.0.1:6379/0') # 创建任务函数@app.taskdef my_task(): print('任务正在执行...')
celery第一个参数是给其设定一个名字, 第二参数我们设定一个中间人broker, 在这里我们使用redis作为中间人。my_task函数是我们编写的一个任务函数, 通过加上装饰器app.task, 将其注册到broker的队列中。
2. 启动redis、创建worker
现在我们在创建一个worker, 等待处理队列中的任务。
进入项目的根目录,执行命令: celery -a celery_tasks.tasks worker -l info
3. 调用任务
下面来测试一下功能,创建一个任务,加入任务队列中,提供worker执行。
进入python终端, 执行如下代码:
$ python manage.py shell>>> from blog.tasks import my_task>>> my_task.delay()<asyncresult: 83484dfe-f729-417b-8e51-6c7ae32a1377>
调用一个任务函数,将会返回一个asyncresult对象,这个对象可以用来检查任务的状态或者获得任务的返回值。
4. 查看结果
在worker的终端查看任务执行情况,可以看到已经收到83484dfe-f729-417b-8e51-6c7ae32a1377 任务,并打印了任务执行信息
5. 存储并查看任务执行状态
把任务执行结果赋值给ret,然后调用result() 会产生 disabledbackend 报错,可见没有配置后端存储的时候并不能保存任务执行的状态信息,下一节我们会讲到如何配置backend保存任务执行结果
$ python manage.py shell>>> from blog.tasks import my_task>>> ret=my_task.delay()>>> ret.result()
四、配置backend存储任务执行结果 如果我们想跟踪任务的状态,celery需要将结果保存到某个地方。有几种保存的方案可选:sqlalchemy、django orm、memcached、 redis、rpc (rabbitmq/amqp)。
1. 添加backend参数
在本例中我们使用redis作为存储结果的方案,通过celery的backend参数来设定任务结果存储地址。我们将tasks模块修改如下:
from celery import celery # 使用redis作为broker以及backendapp = celery('celery_tasks.tasks', broker='redis://127.0.0.1:6379/8', backend='redis://127.0.0.1:6379/9') # 创建任务函数@app.taskdef my_task(a, b): print("任务函数正在执行....") return a + b
给celery增加了backend参数,指定redis作为结果存储,并将任务函数修改为两个参数,并且有返回值。
2. 调用任务/查看任务执行结果
下面再来执行调用一下这个任务看看。
$ python manage.py shell>>> from blog.tasks import my_task>>> res=my_task.delay(10,40)>>> res.result50>>> res.failed()false
再来看看worker的执行情况,如下:
可以看到celery任务已经执行成功了。
但是这只是一个开始,下一步要看看如何添加定时的任务。
四、优化celery目录结构上面直接将celery的应用创建、配置、tasks任务全部写在了一个文件,这样在后面项目越来越大,也是不方便的。下面来拆分一下,并且添加一些常用的参数。
基本结构如下
$ vim typeidea/celery.py (celery应用文件)
#!/usr/bin/env python# -*- coding: utf-8 -*- """#file: celery.py#time: 2022/3/30 12:25 下午#author: julius"""import osfrom celery import celeryfrom blog import celeryconfigproject_name='typeidea'# set the default django setting module for the 'celery' programos.environ.setdefault('django_settings_module','typeidea.settings')app = celery(project_name) app.config_from_object('django.conf:settings') app.autodiscover_tasks()
vim blog/celeryconfig.py (配置celery的参数文件)
#!/usr/bin/env python# -*- coding: utf-8 -*- """#file: celeryconfig.py#time: 2022/3/30 2:54 下午#author: julius"""# 设置结果存储from typeidea import settingsimport os os.environ.setdefault("django_settings_module", "typeidea.settings")celery_result_backend = 'redis://127.0.0.1:6379/0'# 设置代理人brokerbroker_url = 'redis://127.0.0.1:6379/1'# celery 的启动工作数量设置celery_worker_concurrency = 20# 任务预取功能,就是每个工作的进程/线程在获取任务的时候,会尽量多拿 n 个,以保证获取的通讯成本可以压缩。celeryd_prefetch_multiplier = 20# 非常重要,有些情况下可以防止死锁celeryd_force_execv = true# celery 的 worker 执行多少个任务后进行重启操作celery_worker_max_tasks_per_child = 100# 禁用所有速度限制,如果网络资源有限,不建议开足马力。celery_disable_rate_limits = true celery_enable_utc = falsecelery_timezone = settings.time_zonedjango_celery_beat_tz_aware = falsecelery_beat_scheduler = 'django_celery_beat.schedulers:databasescheduler'
vim blog/tasks.py (tasks 任务文件)
import timefrom blog.celery import app # 创建任务函数@app.taskdef my_task(a, b, c): print('任务正在执行...') print('任务1函数休眠10s') time.sleep(10) return a + b + c
五、开始使用django-celery-beat调度器使用 django-celery-beat 动态添加定时任务 celery 4.x 版本在 django 框架中是使用 django-celery-beat 进行动态添加定时任务的。前面虽然已经安装了这个库,但是还要再说明一下。
1. 安装 django-celery-beat
pip3 install django-celery-beat
2.在项目的 settings 文件配置 django-celery-beat
installed_apps = [ 'blog', 'django_celery_beat', ...] # django设置时区language_code = 'zh-hans' # 使用中国语言time_zone = 'asia/shanghai' # 设置django使用中国上海时间# 如果use_tz设置为true时,django会使用系统默认设置的时区,此时的time_zone不管有没有设置都不起作用# 如果use_tz 设置为false,time_zone = 'asia/shanghai', 则使用上海的utc时间。use_tz = false
3. 创建 django-celery-beat 相关表
执行django数据库迁移: python manage.py migrate
4. 配置celery使用 django-celery-beat
配置 celery.py
import os from celery import celery from blog import celeryconfig # 为celery 设置环境变量os.environ.setdefault("django_settings_module","typeidea.settings")# 创建celery appapp = celery('blog')# 从单独的配置模块中加载配置app.config_from_object(celeryconfig) # 设置app自动加载任务app.autodiscover_tasks([ 'blog',])
配置 celeryconfig.py
# 设置结果存储from typeidea import settingsimport os os.environ.setdefault("django_settings_module", "typeidea.settings")celery_result_backend = 'redis://127.0.0.1:6379/0'# 设置代理人brokerbroker_url = 'redis://127.0.0.1:6379/1'# celery 的启动工作数量设置celery_worker_concurrency = 20# 任务预取功能,就是每个工作的进程/线程在获取任务的时候,会尽量多拿 n 个,以保证获取的通讯成本可以压缩。celeryd_prefetch_multiplier = 20# 非常重要,有些情况下可以防止死锁celeryd_force_execv = true# celery 的 worker 执行多少个任务后进行重启操作celery_worker_max_tasks_per_child = 100# 禁用所有速度限制,如果网络资源有限,不建议开足马力。celery_disable_rate_limits = true celery_enable_utc = falsecelery_timezone = settings.time_zonedjango_celery_beat_tz_aware = falsecelery_beat_scheduler = 'django_celery_beat.schedulers:databasescheduler'
编写任务 tasks.py
import timefrom celery import celeryfrom blog.celery import app # 使用redis做为broker# app = celery('blog.tasks2',broker='redis://127.0.0.1:6379/0',backend='redis://127.0.0.1:6379/1') # 创建任务函数@app.taskdef my_task(a, b, c): print('任务正在执行...') print('任务1函数休眠10s') time.sleep(10) return a + b + c @app.taskdef my_task2(): print("任务2函数正在执行....") print('任务2函数休眠10s') time.sleep(10)
5. 启动定时任务work
启动定时任务首先需要有一个work执行异步任务,然后再启动一个定时器触发任务。
启动任务 work
$ celery -a blog worker -l info
启动定时器触发 beat
celery -a blog beat -l info --scheduler django_celery_beat.schedulers:databasescheduler
六、具体操作演练6.1 创建基于间隔时间的周期性任务1. 初始化周期间隔对象interval 对象
>>> from django_celery_beat.models import periodictask, intervalschedule>>> schedule, created = intervalschedule.objects.get_or_create( ... every=10, ... period=intervalschedule.seconds, ... )>>> intervalschedule.objects.all()<queryset [<intervalschedule: every 10 seconds>]>
2.创建一个无参数的周期性间隔任务
>>>periodictask.objects.create(interval=schedule,name='my_task2',task='blog.tasks.my_task2',)<periodictask: my_task2: every 10 seconds>
beat 调度服务日志显示如下:
worker 服务日志显示如下:
3.创建一个带参数的周期性间隔任务
>>> periodictask.objects.create(interval=schedule,name='my_task',task='blog.tasks.my_task',args=json.dumps([10,20,30]))<periodictask: my_task: every 10 seconds>
beat 调度服务日志结果:
worker 服务日志结果:
4.如何高并发执行任务
需要并行执行任务的时候,就需要设置多个worker来执行任务。
6.2 创建一个不带参数的周期性间隔任务1.初始化 crontab 的调度对象
>>> import pytz>>> schedule, _ = crontabschedule.objects.get_or_create(... minute='*',... hour='*',... day_of_week='*',... day_of_month='*',... timezone=pytz.timezone('asia/shanghai')... )
2. 创建不带参数的定时任务
periodictask.objects.create(crontab=schedule,name='my_task2_crontab',task='blog.tasks.my_task2',)
beat 调度服务执行结果
worker 执行服务结果
6.3 周期性任务的查询、删除操作1. 周期性任务的查询
>>> periodictask.objects.all()<extendedqueryset [<periodictask: celery.backend_cleanup: 0 4 * * * (m/h/dm/my/d) asia/shanghai>, <periodictask: my_task2_crontab: * * * * * (m/h/dm/my/d) asia/shanghai>]>>>> periodictask.objects.get(name='my_task2_crontab')<periodictask: my_task2_crontab: * * * * * (m/h/dm/my/d) asia/shanghai>>>> for task in periodictask.objects.all():... print(task.id)... 113>>> periodictask.objects.get(id=13)<periodictask: my_task2_crontab: * * * * * (m/h/dm/my/d) asia/shanghai>>>> periodictask.objects.get(name='my_task2_crontab')<periodictask: my_task2_crontab: * * * * * (m/h/dm/my/d) asia/shanghai>
控制台实际操作记录
2.周期性任务的暂停/启动
2.1 设置my_taks2_crontab 暂停任务
>>> my_task2_crontab = periodictask.objects.get(id=13)>>> my_task2_crontab.enabledtrue>>> my_task2_crontab.enabled=false>>> my_task2_crontab.save()
查看worker输出:
可以看到worker从19:31以后已经没有输出了,说明已经成功吧my_task2_crontab 任务暂停
2.2 设置my_task2_crontab 开启任务
把任务的 enabled 为 true 即可:
>>> my_task2_crontab.enabledfalse>>> my_task2_crontab.enabled=true>>> my_task2_crontab.save()
查看worker输出:
可以看到worker从19:36开始有输出,说明已把my_task2_crontab 任务重新启动
3. 周期性任务的删除
获取到指定的任务后调用delete(),再次查询指定任务会发现已经不存在了
periodictask.objects.get(name='my_task2_crontab').delete()>>> periodictask.objects.get(name='my_task2_crontab')traceback (most recent call last): file "<console>", line 1, in <module> file "/users/julius/pycharmprojects/typeidea/.venv/lib/python3.9/site-packages/django/db/models/manager.py", line 85, in manager_method return getattr(self.get_queryset(), name)(*args, **kwargs) file "/users/julius/pycharmprojects/typeidea/.venv/lib/python3.9/site-packages/django/db/models/query.py", line 435, in get raise self.model.doesnotexist(django_celery_beat.models.periodictask.doesnotexist: periodictask matching query does not exist.
以上就是怎么用python celery动态添加定时任务的详细内容。
