Celery
1.什么是Clelery
Celery是一个简单、灵活且可靠的,处理大量消息的分布式系统
专注于实时处理的异步任务队列
同时也支持任务调度
Celery架构
Celery的架构由三部分组成,消息中间件(message broker),任务执行单元(worker)和任务执行结果存储(task result store)组成。
消息中间件
Celery本身不提供消息服务,但是可以方便的和第三方提供的消息中间件集成。包括,RabbitMQ, Redis等等
任务执行单元
Worker是Celery提供的任务执行的单元,worker并发的运行在分布式的系统节点中。
任务结果存储
Task result store用来存储Worker执行的任务的结果,Celery支持以不同方式存储任务的结果,包括AMQP, redis等
版本支持情况
1Celery version 4.0 runs on 2 Python ❨2.7, 3.4, 3.5❩ 3 PyPy ❨5.4, 5.5❩ 4 This is the last version to support Python 2.7, and from the next version (Celery 5.x) Python 3.5 or newer is required. 5 6 If you’re running an older version of Python, you need to be running an older version of Celery: 7 8 Python 2.6: Celery series 3.1 or earlier. 9 Python 2.5: Celery series 3.0 or earlier. 10 Python 2.4 was Celery series 2.2 or earlier. 11 12 Celery is a project with minimal funding, so we don’t support Microsoft Windows. Please don’t open any issues related to that platform.
2.使用场景
异步任务:将耗时操作任务提交给Celery去异步执行,比如发送短信/邮件、消息推送、音视频处理等等
定时任务:定时执行某件事情,比如每天数据统计
3.Celery的安装配置
pip install celery
消息中间件:RabbitMQ/Redis
app=Celery('任务名',backend='xxx',broker='xxx')
4.Celery执行异步任务
基本使用
创建项目celerytest
创建py文件:celery_app_task.py
1import celery 2import time 3# broker='redis://127.0.0.1:6379/2' 不加密码 4backend='redis://:123456@127.0.0.1:6379/1' 5broker='redis://:123456@127.0.0.1:6379/2' 6cel=celery.Celery('test',backend=backend,broker=broker) 7@cel.task 8def add(x,y): 9 return x+y
创建py文件:add_task.py,添加任务
1from celery_app_task import add 2result = add.delay(4,5) 3print(result.id)
创建py文件:run.py,执行任务,或者使用命令执行:celery worker -A celery_app_task -l info
注:windows下:celery worker -A celery_app_task -l info -P eventlet
1from celery_app_task import cel 2if __name__ == '__main__': 3 cel.worker_main() 4 # cel.worker_main(argv=['--loglevel=info')
创建py文件:result.py,查看任务执行结果
1from celery.result import AsyncResult 2from celery_app_task import cel 3 4async = AsyncResult(id="e919d97d-2938-4d0f-9265-fd8237dc2aa3", app=cel) 5 6if async.successful(): 7 result = async.get() 8 print(result) 9 # result.forget() # 将结果删除 10elif async.failed(): 11 print('执行失败') 12elif async.status == 'PENDING': 13 print('任务等待中被执行') 14elif async.status == 'RETRY': 15 print('任务异常后正在重试') 16elif async.status == 'STARTED': 17 print('任务已经开始被执行')
执行 add_task.py,添加任务,并获取任务ID
执行 run.py ,或者执行命令:celery worker -A celery_app_task -l info
执行 result.py,检查任务状态并获取结果
多任务结构
1pro_cel 2 ├── celery_task# celery相关文件夹 3 │ ├── celery.py # celery连接和配置相关文件,必须叫这个名字 4 │ └── tasks1.py # 所有任务函数 5 │ └── tasks2.py # 所有任务函数 6 ├── check_result.py # 检查结果 7 └── send_task.py # 触发任务
celery.py
1from celery import Celery 2 3cel = Celery('celery_demo', 4 broker='redis://127.0.0.1:6379/1', 5 backend='redis://127.0.0.1:6379/2', 6 # 包含以下两个任务文件,去相应的py文件中找任务,对多个任务做分类 7 include=['celery_task.tasks1', 8 'celery_task.tasks2' 9 ]) 10 11# 时区 12cel.conf.timezone = 'Asia/Shanghai' 13# 是否使用UTC 14cel.conf.enable_utc = False
tasks1.py
1import time 2from celery_task.celery import cel 3 4@cel.task 5def test_celery(res): 6 time.sleep(5) 7 return "test_celery任务结果:%s"%res
tasks2.py
1import time 2from celery_task.celery import cel 3@cel.task 4def test_celery2(res): 5 time.sleep(5) 6 return "test_celery2任务结果:%s"%res
check_result.py
1from celery.result import AsyncResult 2from celery_task.celery import cel 3 4async = AsyncResult(id="08eb2778-24e1-44e4-a54b-56990b3519ef", app=cel) 5 6if async.successful(): 7 result = async.get() 8 print(result) 9 # result.forget() # 将结果删除,执行完成,结果不会自动删除 10 # async.revoke(terminate=True) # 无论现在是什么时候,都要终止 11 # async.revoke(terminate=False) # 如果任务还没有开始执行呢,那么就可以终止。 12elif async.failed(): 13 print('执行失败') 14elif async.status == 'PENDING': 15 print('任务等待中被执行') 16elif async.status == 'RETRY': 17 print('任务异常后正在重试') 18elif async.status == 'STARTED': 19 print('任务已经开始被执行')
send_task.py
1from celery_task.tasks1 import test_celery 2from celery_task.tasks2 import test_celery2 3 4# 立即告知celery去执行test_celery任务,并传入一个参数 5result = test_celery.delay('第一个的执行') 6print(result.id) 7result = test_celery2.delay('第二个的执行') 8print(result.id)
添加任务(执行send_task.py),开启work:celery worker -A celery_task -l info -P eventlet,检查任务执行结果(执行check_result.py)
5.Celery执行定时任务
设定时间让celery执行一个任务
add_task.py
1from celery_app_task import add 2from datetime import datetime 3 4# 方式一 5# v1 = datetime(2019, 2, 13, 18, 19, 56) 6# print(v1) 7# v2 = datetime.utcfromtimestamp(v1.timestamp()) 8# print(v2) 9# result = add.apply_async(args=[1, 3], eta=v2) 10# print(result.id) 11 12# 方式二 13ctime = datetime.now() 14# 默认用utc时间 15utc_ctime = datetime.utcfromtimestamp(ctime.timestamp()) 16from datetime import timedelta 17time_delay = timedelta(seconds=10) 18task_time = utc_ctime + time_delay 19 20# 使用apply_async并设定时间 21result = add.apply_async(args=[4, 3], eta=task_time) 22print(result.id)
类似于contab的定时任务
多任务结构中celery.py修改如下
1from datetime import timedelta 2from celery import Celery 3from celery.schedules import crontab 4 5cel = Celery('tasks', broker='redis://127.0.0.1:6379/1', backend='redis://127.0.0.1:6379/2', include=[ 6 'celery_task.tasks1', 7 'celery_task.tasks2', 8]) 9cel.conf.timezone = 'Asia/Shanghai' 10cel.conf.enable_utc = False 11 12cel.conf.beat_schedule = { 13 # 名字随意命名 14 'add-every-10-seconds': { 15 # 执行tasks1下的test_celery函数 16 'task': 'celery_task.tasks1.test_celery', 17 # 每隔2秒执行一次 18 # 'schedule': 1.0, 19 # 'schedule': crontab(minute="*/1"), 20 'schedule': timedelta(seconds=2), 21 # 传递参数 22 'args': ('test',) 23 }, 24 # 'add-every-12-seconds': { 25 # 'task': 'celery_task.tasks1.test_celery', 26 # 每年4月11号,8点42分执行 27 # 'schedule': crontab(minute=42, hour=8, day_of_month=11, month_of_year=4), 28 # 'schedule': crontab(minute=42, hour=8, day_of_month=11, month_of_year=4), 29 # 'args': (16, 16) 30 # }, 31}
启动一个beat:celery beat -A celery_task -l info
启动work执行:celery worker -A celery_task -l info -P eventlet
6.Django中使用Celery
在项目目录下创建celeryconfig.py
1import djcelery 2djcelery.setup_loader() 3CELERY_IMPORTS=( 4 'app01.tasks', 5) 6#有些情况可以防止死锁 7CELERYD_FORCE_EXECV=True 8# 设置并发worker数量 9CELERYD_CONCURRENCY=4 10#允许重试 11CELERY_ACKS_LATE=True 12# 每个worker最多执行100个任务被销毁,可以防止内存泄漏 13CELERYD_MAX_TASKS_PER_CHILD=100 14# 超时时间 15CELERYD_TASK_TIME_LIMIT=12*30
在app01目录下创建tasks.py
1from celery import task 2@task 3def add(a,b): 4 with open('a.text', 'a', encoding='utf-8') as f: 5 f.write('a') 6 print(a+b)
视图函数views.py
1from django.shortcuts import render,HttpResponse 2from app01.tasks import add 3from datetime import datetime 4def test(request): 5 # result=add.delay(2,3) 6 ctime = datetime.now() 7 # 默认用utc时间 8 utc_ctime = datetime.utcfromtimestamp(ctime.timestamp()) 9 from datetime import timedelta 10 time_delay = timedelta(seconds=5) 11 task_time = utc_ctime + time_delay 12 result = add.apply_async(args=[4, 3], eta=task_time) 13 print(result.id) 14 return HttpResponse('ok')
settings.py
1INSTALLED_APPS = [ 2 ... 3 'djcelery', 4 'app01' 5] 6 7... 8 9from djagocele import celeryconfig 10BROKER_BACKEND='redis' 11BOOKER_URL='redis://127.0.0.1:6379/1' 12CELERY_RESULT_BACKEND='redis://127.0.0.1:6379/2'