欢迎您访问程序员文章站本站旨在为大家提供分享程序员计算机编程知识!
您现在的位置是: 首页  >  IT编程

Python 并行分布式框架 Celery

程序员文章站 2024-01-05 10:10:28
Python 并行分布式框架 Celery Python 并行分布式框架 Celery Python 并行分布式框架 Celery Python 并行分布式框架 Celery Celery 官网:http://www.celeryproject.orgCelery 官方文档英文版:http://do ......

python 并行分布式框架 celery

celery 官网:
celery 官方文档英文版:
celery 官方文档中文版:

celery配置:

参考:    http://blog.csdn.net/happyanger6/article/details/51408266

分布式队列神器 celery:
celery最佳实践:
celery 分布式任务队列快速入门:
异步任务神器 celery 快速入门教程:
定时任务管理之python篇celery使用:
异步任务神器 celery:
celery任务调度框架实践:
celery-4.1 用户指南: monitoring and management guide:
celery安装及使用:
celery学习笔记(一):

 

 

 

celery 简介

 

        除了redis,还可以使用另外一个神器---celery。celery是一个异步任务的调度工具。

        celery 是 distributed task queue,分布式任务队列,分布式决定了可以有多个 worker 的存在,队列表示其是异步操作,即存在一个产生任务提出需求的工头,和一群等着被分配工作的码农。

        在 python 中定义 celery 的时候,我们要引入 broker,中文翻译过来就是“中间人”的意思,在这里 broker 起到一个中间人的角色。在工头提出任务的时候,把所有的任务放到 broker 里面,在 broker 的另外一头,一群码农等着取出一个个任务准备着手做。

        这种模式注定了整个系统会是个开环系统,工头对于码农们把任务做的怎样是不知情的。所以我们要引入 backend 来保存每次任务的结果。这个 backend 有点像我们的 broker,也是存储任务的信息用的,只不过这里存的是那些任务的返回结果。我们可以选择只让错误执行的任务返回结果到 backend,这样我们取回结果,便可以知道有多少任务执行失败了。

        celery(芹菜)是一个异步任务队列/基于分布式消息传递的作业队列。它侧重于实时操作,但对调度支持也很好。celery用于生产系统每天处理数以百万计的任务。celery是用python编写的,但该协议可以在任何语言实现。它也可以与其他语言通过webhooks实现。celery建议的消息队列是rabbitmq,但提供有限支持redis, beanstalk, mongodb, couchdb, 和数据库(使用sqlalchemy的或django的 orm) 。celery是易于集成django, pylons and flask,使用 django-celery, celery-pylons and flask-celery 附加包即可。

 

 

 

celery 介绍

 

在celery中几个基本的概念,需要先了解下,不然不知道为什么要安装下面的东西。概念:broker、backend。

什么是broker?

broker是一个消息传输的中间件,可以理解为一个邮箱。每当应用程序调用celery的异步任务的时候,会向broker传递消息,而后celery的worker将会取到消息,进行对于的程序执行。好吧,这个邮箱可以看成是一个消息队列。其中broker的中文意思是 经纪人 ,其实就是一开始说的 消息队列 ,用来发送和接受消息。这个broker有几个方案可供选择:rabbitmq (消息队列),(缓存数据库),(不推荐),等等

什么是backend?

通常程序发送的消息,发完就完了,可能都不知道对方时候接受了。为此,celery实现了一个backend,用于存储这些消息以及celery执行的一些消息和结果。backend是在celery的配置中的一个配置项 celery_result_backend ,作用是保存结果和状态,如果你需要跟踪任务的状态,那么需要设置这一项,可以是database backend,也可以是cache backend,具体可以参考这里: celery_result_backend 。

对于 brokers,官方推荐是 rabbitmq 和 redis,至于 backend,就是数据库。为了简单可以都使用 redis。

我自己演示使用rabbitmq作为broker,用作为backend。

来一张图,这是在网上最多的一张celery的图了,确实描述的非常好

Python 并行分布式框架 Celery

celery的由三部分组成,消息中间件(message broker),任务执行单元(worker)和任务执行结果存储(task result store)组成。

消息中间件

celery本身不提供消息服务,但是可以方便的和第三方提供的消息中间件集成。包括,rabbitmq, ,  (experimental), amazon sqs (experimental),couchdb (experimental), sqlalchemy (experimental),django orm (experimental), ironmq

任务执行单元

worker是celery提供的任务执行的单元,worker并发的运行在分布式的系统节点中。

任务结果存储

task result store用来存储worker执行的任务的结果,celery支持以不同方式存储任务的结果,包括amqp, ,memcached, ,sqlalchemy, django orm,apache cassandra, ironcache 等。

这里我先不去看它是如何存储的,就先选用redis来存储任务执行结果。

因为涉及到消息中间件(在celery帮助文档中称呼为中间人<broker>),为了更好的去理解文档中的例子,可以安装两个中间件,一个是rabbitmq,一个redis。

根据 celery的帮助文档 安装和设置rabbitmq, 要使用 celery,需要创建一个 rabbitmq 用户、一个虚拟主机,并且允许这个用户访问这个虚拟主机。

$ sudo rabbitmqctl add_user forward password     #创建了一个rabbitmq用户,用户名为forward,密码是password
$ sudo rabbitmqctl add_vhost ubuntu              #创建了一个虚拟主机,主机名为ubuntu

# 设置权限。允许用户forward访问虚拟主机ubuntu,因为rabbitmq通过主机名来与节点通信
$ sudo rabbitmqctl set_permissions -p ubuntu forward ".*" ".*" ".*"     
$ sudo rabbitmq-server    # 启用rabbitmq服务器

结果如下,成功运行:

Python 并行分布式框架 Celery

安装redis,它的安装比较简单

$ sudo pip install redis

然后进行简单的配置,只需要设置 redis 数据库的位置:
broker_url = 'redis://localhost:6379/0'

url的格式为:
redis://:password@hostname:port/db_number
url scheme 后的所有字段都是可选的,并且默认为 localhost 的 6479 端口,使用数据库 0。我的配置是:

redis://:password@ubuntu:6379/5

安装celery,我是用标准的python工具pip安装的,如下:

$ sudo pip install celery

 

 

celery 是一个强大的 分布式任务队列 的 异步处理框架,它可以让任务的执行完全脱离主程序,甚至可以被分配到其他主机上运行。我们通常使用它来实现异步任务(async task)和定时任务(crontab)。我们需要一个消息队列来下发我们的任务。首先要有一个消息中间件,此处选择rabbitmq (也可选择 redis 或 amazon simple queue service(sqs)消息队列服务)。推荐 选择 rabbitmq 。使用rabbitmq是官方特别推荐的方式,因此我也使用它作为我们的broker。它的组成如下图:

Python 并行分布式框架 Celery

 

可以看到,celery 主要包含以下几个模块:

  • 任务模块 task

    包含异步任务和定时任务。其中,异步任务通常在业务逻辑中被触发并发往任务队列,而定时任务由 celery beat 进程周期性地将任务发往任务队列。

  • 消息中间件 broker

    broker,即为任务调度队列,接收任务生产者发来的消息(即任务),将任务存入队列。celery 本身不提供队列服务,官方推荐使用 rabbitmq 和  等。

  • 任务执行单元 worker

    worker 是执行任务的处理单元,它实时监控消息队列,获取队列中调度的任务,并执行它。

  • 任务结果存储 backend

    backend 用于存储任务的执行结果,以供查询。同消息中间件一样,存储也可使用 rabbitmq,  和  等。

 

 

 

 

安装

 

有了上面的概念,需要安装这么几个东西:rabbitmq、sqlalchemy、celery

安装rabbitmq

官网安装方法:

启动管理插件:sbin/rabbitmq-plugins enable rabbitmq_management 
启动rabbitmq:sbin/rabbitmq-server -detached

rabbitmq已经启动,可以打开页面来看看 
地址: 

用户名密码都是guest 。现在可以进来了,可以看到具体页面。 关于rabbitmq的配置,网上很多 自己去搜以下就ok了。

消息中间件有了,现在该来代码了,使用  celeby官网代码。

剩下两个都是的东西了,直接pip安装就好了,对于从来没有安装过驱动的同学可能需要安装mysql-。安装完成之后,启动服务: $ rabbitmq-server[回车]。启动后不要关闭窗口, 下面操作新建窗口(tab)。

 

安装celery
celery可以通过pip自动安装,如果你喜欢使用虚拟环境安装可以先使用virtualenv创建一个自己的虚拟环境。反正我喜欢使用virtualenv建立自己的环境。

pip install celery

 

http://www.open-open.com/lib/view/open1441161168878.html

 

开始使用 celery

 

使用celery包含三个方面:1. 定义任务函数。2. 运行celery服务。3. 客户应用程序的调用。

创建一个文件 tasks.py输入下列代码:

  1.  
    from celery import celery
  2.  
     
  3.  
    broker = 'redis://127.0.0.1:6379/5'
  4.  
    backend = 'redis://127.0.0.1:6379/6'
  5.  
     
  6.  
     
  7.  
    app = celery('tasks', broker=broker, backend=backend)
  8.  
     
  9.  
    @app.task
  10.  
    def add(x, y):
  11.  
    return x + y

上述代码导入了celery,然后创建了celery 实例 app,实例化的过程中指定了任务名tasks(和文件名一致),传入了broker和backend。然后创建了一个任务函数add。下面启动celery服务。在当前命令行终端运行(分别在 env1 和 env2 下执行):

celery -a tasks worker  --loglevel=info

目录结构 (celery -a tasks worker --loglevel=info 这条命令当前工作目录必须和 tasks.py 所在的目录相同。即 进入tasks.py所在目录执行这条命令。

Python 并行分布式框架 Celery

使用 python 虚拟环境 模拟两个不同的 主机。

Python 并行分布式框架 Celery

此时会看见一对输出。包括注册的任务啦。

Python 并行分布式框架 Celery

 

交互式客户端程序调用方法

打开一个命令行,进入python环境。

in [0]:from tasks import add
in [1]: r = add.delay(2, 2)
in [2]: add.delay(2, 2)
out[2]: <asyncresult: 6fdb0629-4beb-4eb7-be47-f22be1395e1d>

in [3]: r = add.delay(3, 3)

in [4]: r.re
r.ready   r.result  r.revoke

in [4]: r.ready()
out[4]: true

in [6]: r.result
out[6]: 6

in [7]: r.get()
out[7]: 6

Python 并行分布式框架 Celery

调用 delay 函数即可启动 add 这个任务。这个函数的效果是发送一条消息到broker中去,这个消息包括要执行的函数、函数的参数以及其他信息,具体的可以看 celery官方文档。这个时候 worker 会等待 broker 中的消息,一旦收到消息就会立刻执行消息。

启动了一个任务之后,可以看到之前启动的worker已经开始执行任务了。

Python 并行分布式框架 Celery

现在是在python环境中调用的add函数,实际上通常在应用程序中调用这个方法。

注意:如果把返回值赋值给一个变量,那么原来的应用程序也会被阻塞,需要等待异步任务返回的结果。因此,实际使用中,不需要把结果赋值。

 

应用程序中调用方法

新建一个 main.py 文件 代码如下:

  1.  
    from tasks import add
  2.  
     
  3.  
    r = add.delay(2, 2)
  4.  
    r = add.delay(3, 3)
  5.  
    print r.ready()
  6.  
    print r.result
  7.  
    print r.get()

Python 并行分布式框架 Celery

在celery命令行可以看见celery执行的日志。打开 backend的redis,也可以看见celery执行的信息。

使用  redis desktop manager 查看 redis 数据库内容如图:

Python 并行分布式框架 Celery

 

使用配置文件

celery 的配置比较多,可以在 官方配置文档:  查询每个配置项的含义。

上述的使用是简单的配置,下面介绍一个更健壮的方式来使用celery。首先创建一个python包,celery服务,姑且命名为proj。目录文件如下:

  1.  
     
  2.  
    proj tree
  3.  
    .
  4.  
    ├── __init__.py
  5.  
    ├── celery.py # 创建 celery 实例
  6.  
    ├── config.py # 配置文件
  7.  
    └── tasks.py # 任务函数

首先是 celery.py

  1.  
    #!/usr/bin/env python
  2.  
    # -*- coding:utf-8 -*-
  3.  
     
  4.  
    from __future__ import absolute_import
  5.  
    from celery import celery
  6.  
     
  7.  
    app = celery('proj', include=['proj.tasks'])
  8.  
     
  9.  
    app.config_from_object('proj.config')
  10.  
     
  11.  
    if __name__ == '__main__':
  12.  
    app.start()

这一次创建 app,并没有直接指定 broker 和 backend。而是在配置文件中。

config.py

  1.  
    #!/usr/bin/env python
  2.  
    # -*- coding:utf-8 -*-
  3.  
     
  4.  
    from __future__ import absolute_import
  5.  
     
  6.  
    celery_result_backend = 'redis://127.0.0.1:6379/5'
  7.  
    broker_url = 'redis://127.0.0.1:6379/6'

剩下的就是tasks.py

  1.  
    #!/usr/bin/env python
  2.  
    # -*- coding:utf-8 -*-
  3.  
     
  4.  
    from __future__ import absolute_import
  5.  
    from proj.celery import app
  6.  
     
  7.  
    @app.task
  8.  
    def add(x, y):
  9.  
    return x + y

使用方法也很简单,在 proj 的同一级目录执行 celery

celery -a proj worker -l info

现在使用任务也很简单,直接在客户端代码调用 proj.tasks 里的函数即可。

 

指定 路由 到的 队列

celery的官方文档 。先看代码(tasks.py):

  1.  
    from celery import celery
  2.  
     
  3.  
    app = celery()
  4.  
    app.config_from_object("celeryconfig")
  5.  
     
  6.  
    @app.task
  7.  
    def taska(x,y):
  8.  
    return x + y
  9.  
     
  10.  
    @app.task
  11.  
    def taskb(x,y,z):
  12.  
    return x + y + z
  13.  
     
  14.  
    @app.task
  15.  
    def add(x,y):
  16.  
    return x + y

上面的tasks.py中,首先定义了一个celery对象,然后用celeryconfig.py对celery对象进行设置,之后再分别定义了三个task,分别是taska,taskb和add。接下来看一下celeryconfig.py 文件

  1.  
    from kombu import exchange,queue
  2.  
     
  3.  
    broker_url = "redis://10.32.105.227:6379/0" celery_result_backend = "redis://10.32.105.227:6379/0"
  4.  
     
  5.  
    celery_queues = (
  6.  
       queue("default",exchange("default"),routing_key="default"),
  7.  
       queue("for_task_a",exchange("for_task_a"),routing_key="task_a"),
  8.  
       queue("for_task_b",exchange("for_task_b"),routing_key="task_a")
  9.  
     )
  10.  
     
  11.  
    celery_routes = {
  12.  
    'tasks.taska':{"queue":"for_task_a","routing_key":"task_a"},
  13.  
    'tasks.taskb":{"queue":"for_task_b","routing_key:"task_b"}
  14.  
    }

在 celeryconfig.py 文件中,首先设置了brokel以及result_backend,接下来定义了三个message queue,并且指明了queue对应的exchange(当使用redis作为broker时,exchange的名字必须和queue的名字一样)以及routing_key的值。

现在在一台主机上面启动一个worker,这个worker只执行for_task_a队列中的消息,这是通过在启动worker是使用-q queue_name参数指定的。

celery -a tasks worker -l info -n worker.%h -q for_task_a

然后到另一台主机上面执行taska任务。首先 切换当前目录到代码所在的工程下,启动python,执行下面代码启动taska:

  1.  
    from tasks import *
  2.  
     
  3.  
    task_a_re = taska.delay(100,200)

执行完上面的代码之后,task_a消息会被立即发送到for_task_a队列中去。此时已经启动的worker.atsgxxx 会立即执行taska任务。

重复上面的过程,在另外一台机器上启动一个worker专门执行for_task_b中的任务。修改上一步骤的代码,把 taska 改成 taskb 并执行。

  1.  
    from tasks import *
  2.  
     
  3.  
    task_b_re = taskb.delay(100,200)

在上面的 tasks.py 文件中还定义了add任务,但是在celeryconfig.py文件中没有指定这个任务route到那个queue中去执行,此时执行add任务的时候,add会route到celery默认的名字叫做celery的队列中去。

Python 并行分布式框架 Celery

因为这个消息没有在celeryconfig.py文件中指定应该route到哪一个queue中,所以会被发送到默认的名字为celery的queue中,但是我们还没有启动worker执行celery中的任务。接下来我们在启动一个worker执行celery队列中的任务。

celery -a tasks worker -l info -n worker.%h -q celery 

然后再查看add的状态,会发现状态由pending变成了success。

 

scheduler ( 定时任务,周期性任务 )

一种常见的需求是每隔一段时间执行一个任务。

在celery中执行定时任务非常简单,只需要设置celery对象的celerybeat_schedule属性即可。

配置如下

config.py

  1.  
    #!/usr/bin/env python
  2.  
    # -*- coding:utf-8 -*-
  3.  
     
  4.  
    from __future__ import absolute_import
  5.  
     
  6.  
    celery_result_backend = 'redis://127.0.0.1:6379/5'
  7.  
    broker_url = 'redis://127.0.0.1:6379/6'
  8.  
     
  9.  
    celery_timezone = 'asia/shanghai'
  10.  
     
  11.  
    from datetime import timedelta
  12.  
     
  13.  
    celerybeat_schedule = {
  14.  
    'add-every-30-seconds': {
  15.  
    'task': 'proj.tasks.add',
  16.  
    'schedule': timedelta(seconds=30),
  17.  
    'args': (16, 16)
  18.  
    },
  19.  
    }

注意配置文件需要指定时区。这段代码表示每隔30秒执行 add 函数。一旦使用了 scheduler, 启动 celery需要加上-b 参数。

celery -a proj worker -b -l info

设置多个定时任务

celery_timezone = 'utc'
celerybeat_schedule = {
    'taska_schedule' : {
        'task':'tasks.taska',
        'schedule':20,
        'args':(5,6)
    },
    'taskb_scheduler' : {
        'task':"tasks.taskb",
        "schedule":200,
        "args":(10,20,30)
    },
    'add_schedule': {
        "task":"tasks.add",
        "schedule":10,
        "args":(1,2)
    }
}

定义3个定时任务,即每隔20s执行taska任务,参数为(5,6),每隔200s执行taskb任务,参数为(10,20,30),每隔10s执行add任务,参数为(1,2).通过下列命令启动一个定时任务: celery -a tasks beat。使用 beat 参数即可启动定时任务。

 

crontab

计划任务当然也可以用crontab实现,celery也有crontab模式。修改 config.py

  1.  
    #!/usr/bin/env python
  2.  
    # -*- coding:utf-8 -*-
  3.  
     
  4.  
    from __future__ import absolute_import
  5.  
     
  6.  
    celery_result_backend = 'redis://127.0.0.1:6379/5'
  7.  
    broker_url = 'redis://127.0.0.1:6379/6'
  8.  
     
  9.  
    celery_timezone = 'asia/shanghai'
  10.  
     
  11.  
    from celery.schedules import crontab
  12.  
     
  13.  
    celerybeat_schedule = {
  14.  
    # executes every monday morning at 7:30 a.m
  15.  
    'add-every-monday-morning': {
  16.  
    'task': 'tasks.add',
  17.  
    'schedule': crontab(hour=7, minute=30, day_of_week=1),
  18.  
    'args': (16, 16),
  19.  
    },
  20.  
    }

scheduler的切分度很细,可以精确到秒。crontab模式就不用说了。

当然celery还有更高级的用法,比如 多个机器 使用,启用多个 worker并发处理 等。

 

发送任务到队列中

apply_async(args[, kwargs[, …]])、delay(*args, **kwargs) :

send_task  :http://docs.celeryproject.org/en/master/reference/celery.html#celery.celery.send_task

from celery import celery
celery = celery()
celery.config_from_object('celeryconfig')
send_task('tasks.test1', args=[hotplay_id, start_dt, end_dt], queue='hotplay_jy_queue')  

 

 

celery 监控 和 管理  以及 命令帮助

输入 celery -h 可以看到 celery 的命令和帮助

Python 并行分布式框架 Celery

更详细的帮助可以看官方文档:

 

 

celery 官网 示例

 

官网示例:

python 并行分布式框架 celery 超详细介绍

 

 

一个简单例子

 

第一步

编写简单的纯python函数

  1.  
    def say(x,y):
  2.  
    return x+y
  3.  
     
  4.  
    if __name__ == '__main__':
  5.  
    say('hello','world')

 

第二步

如果这个函数不是简单的输出两个字符串相加,而是需要查询数据库或者进行复杂的处理。这种处理需要耗费大量的时间,还是这种方式执行会是多么糟糕的事情。为了演示这种现象,可以使用sleep函数来模拟高耗时任务。

  1.  
    import time
  2.  
     
  3.  
    def say(x,y):
  4.  
    time.sleep(5)
  5.  
    return x+y
  6.  
     
  7.  
    if __name__ == '__main__':
  8.  
    say('hello','world')

 

第三步

这时候我们可能会思考怎么使用多进程或者多线程去实现这种任务。对于多进程与多线程的不足这里不做讨论。现在我们可以想想celery到底能不能解决这种问题。

  1.  
    import time
  2.  
    from celery import celery
  3.  
     
  4.  
    app = celery('sample',broker='amqp://guest@localhost//')
  5.  
     
  6.  
    @app.task
  7.  
    def say(x,y):
  8.  
    time.sleep(5)
  9.  
    return x+y
  10.  
     
  11.  
    if __name__ == '__main__':
  12.  
    say('hello','world')

现在来解释一下新加入的几行代码,首先说明一下加入的新代码完全不需要改变原来的代码。导入celery模块就不用解释了,声明一个celery实例app的参数需要解释一下。

  1. 第一个参数是这个python文件的名字,注意到已经把.py去掉了。
  2. 第二个参数是用到的rabbitmq队列。可以看到其使用的方式非常简单,因为它是默认的消息队列端口号都不需要指明。

 

第四步

现在我们已经使用了celery框架了,我们需要让它找几个工人帮我们干活。好现在就让他们干活。

celery -a sample worker --loglevel=info

这条命令有些长,我来解释一下吧。

  1. -a 代表的是application的首字母,我们的应用就是在 sample 里面 定义的。
  2. worker 就是我们的工人了,他们会努力完成我们的工作的。
  3. -loglevel=info 指明了我们的工作后台执行情况,虽然工人们已经向你保证过一定努力完成任务。但是谨慎的你还是希望看看工作进展情况。
    回车后你可以看到类似下面这样一个输出,如果是没有红色的输出那么你应该是没有遇到什么错误的。

 

第五步

现在我们的任务已经被加载到了内存中,我们不能再像之前那样执行python sample.py来运行程序了。我们可以通过终端进入python然后通过下面的方式加载任务。输入python语句。

  1.  
    from sample import say
  2.  
    say.delay('hello','world')

我们的函数会立即返回,不需要等待。就那么简单celery解决了我们的问题。可以发现我们的say函数不是直接调用了,它被celery 的 task 装饰器 修饰过了。所以多了一些属性。目前我们只需要知道使用delay就行了。

 

简单案例

确保你之前的rabbitmq已经启动。还是官网的那个例子,在任意目录新建一个tasks.py的文件,内容如下:

  1.  
    from celery import celery
  2.  
     
  3.  
    app = celery('tasks', broker='amqp://guest@localhost//')
  4.  
     
  5.  
    @app.task
  6.  
    def add(x, y):
  7.  
    return x + y

使用作为消息队列

  1.  
    app = celery('task', broker='redis://localhost:6379/4')
  2.  
    app.conf.update(
  3.  
    celery_task_serializer='json',
  4.  
    celery_accept_content=['json'], # ignore other content
  5.  
    celery_result_serializer='json',
  6.  
    celeryd_concurrency = 8
  7.  
    )
  8.  
     
  9.  
    @app.task
  10.  
    def add(x, y):
  11.  
    return x + y

在同级目录执行:

$ celery -a tasks.app worker --loglevel=info

该命令的意思是启动一个worker ( tasks文件中的app实例,默认实例名为app,-a 参数后也可直接加文件名,不需要 .app),把tasks中的任务(add(x,y))把任务放到队列中。保持窗口打开,新开一个窗口进入交互模式,python或者ipython:

  1.  
    >>> from tasks import add
  2.  
    >>> add.delay(4, 4)

到此为止,你已经可以使用celery执行任务了,上面的python交互模式下简单的调用了add任务,并传递 4,4 参数。

但此时有一个问题,你突然想知道这个任务的执行结果和状态,到底完了没有。因此就需要设置backend了

修改之前的tasks.py中的代码为:

  1.  
    # coding:utf-8
  2.  
    import subprocess
  3.  
    from time import sleep
  4.  
     
  5.  
    from celery import celery
  6.  
     
  7.  
    backend = 'db+mysql://root:@192.168.0.102/celery'
  8.  
    broker = 'amqp://guest@192.168.0.102:5672'
  9.  
     
  10.  
    app = celery('tasks', backend=backend, broker=broker)
  11.  
     
  12.  
     
  13.  
    @app.task
  14.  
    def add(x, y):
  15.  
    sleep(10)
  16.  
    return x + y
  17.  
     
  18.  
     
  19.  
    @app.task
  20.  
    def hostname():
  21.  
    return subprocess.check_output(['hostname'])

除了添加backend之外,上面还添加了一个who的方法用来多服务器操作。修改完成之后,还按之前的方式启动。

同样进入python的交互模型:

  1.  
    >>> from tasks import add, hostname
  2.  
    >>> r = add.delay(4, 4)
  3.  
    >>> r.ready() # 10s内执行,会输出false,因为add中sleep了10s
  4.  
    >>>
  5.  
    >>> r = hostname.delay()
  6.  
    >>> r.result # 输出你的hostname

 

测试多服务器

做完上面的测试之后,产生了一个疑惑,celery叫做分布式任务管理,那它的分布式体现在哪?它的任务都是怎么执行的?在哪个机器上执行的?在当前服务器上的celery服务不关闭的情况下,按照同样的方式在另外一台服务器上安装celery,并启动:

$ celery -a tasks worker --loglevel=info

发现前一个服务器的celery服务中输出你刚启动的服务器的hostname,前提是那台服务器连上了你的rabbitmq。然后再进入python交互模式:

  1.  
    >>> from tasks import hostname
  2.  
    >>>
  3.  
    >>> for i in range(10):
  4.  
    ... r = hostname.delay()
  5.  
    ... print r.result # 输出你的hostname
  6.  
    >>>

看你输入的内容已经观察两台服务器上你启动celery服务的输出。

 

celery的使用技巧(celery配置文件和发送任务)

在实际的项目中我们需要明确前后台的分界线,因此我们的celery编写的时候就应该是分成前后台两个部分编写。在celery简单入门中的总结部分我们也提出了另外一个问题,就是需要分离celery的配置文件。

第一步

编写后台任务tasks.py脚本文件。在这个文件中我们不需要再声明celery的实例,我们只需要导入其task装饰器来注册我们的任务即可。后台处理业务逻辑完全独立于前台,这里只是简单的hello world程序需要多少个参数只需要告诉前台就可以了,在实际项目中可能你需要的是后台执行发送一封邮件的任务或者进行复杂的数据库查询任务等。

  1.  
    import time
  2.  
    from celery.task import task
  3.  
     
  4.  
     
  5.  
    @task
  6.  
    def say(x,y):
  7.  
    time.sleep(5)
  8.  
    return x+y

第二步

有了那么完美的后台,我们的前台编写肯定也轻松不少。到底简单到什么地步呢,来看看前台的代码吧!为了形象的表明其职能,我们将其命名为client.py脚本文件。

  1.  
    from celery import celery
  2.  
     
  3.  
    app = celery()
  4.  
     
  5.  
    app.config_from_object('celeryconfig')
  6.  
    app.send_task("tasks.say",['hello','world'])

可以看到只需要简单的几步:1. 声明一个celery实例。2. 加载配置文件。3. 发送任务。

第三步

继续完成celery的配置。官方的介绍使用celeryconfig.py作为配置文件名,这样可以防止与你现在的应用的配置同名。

celery_imports = ('tasks')
celery_ignore_result = false
broker_host = '127.0.0.1'
broker_port = 5672
broker_url = 'amqp://'
celery_result_backend = 'amqp'

可以看到我们指定了celery_result_backend为amqp默认的队列!这样我们就可以查看处理后的运行状态了,后面将会介绍处理结果的查看。

第四步

启动celery后台服务,这里是测试与学习celery的教程。在实际生产环境中,如果是通过这种方式启动的后台进程是不行的。所谓后台进程通常是需要作为守护进程运行在后台的,在python的世界里总是有一些工具能够满足你的需要。这里可以使用supervisor作为进程管理工具。在后面的文章中将会介绍如何使用supervisor工具。

celery worker -l info --beat

注意现在运行worker的方式也与前面介绍的不一样了,下面简单介绍各个参数。
    -l info     与--loglevel=info的作用是一样的。
    --beat    周期性的运行。即设置 心跳。

第五步

前台的运行就比较简单了,与平时运行的python脚本一样。python client.py。

现在前台的任务是运行了,可是任务是被写死了。我们的任务大多数时候是动态的,为演示动态工作的情况我们可以使用终端发送任务。

  1.  
    >>> from celery import celery
  2.  
    >>> app = celery()
  3.  
    >>> app.config_from_object('celeryconfig')

在python终端导入celery模块声明实例然后加载配置文件,完成了这些步骤后就可以动态的发送任务并且查看任务状态了。注意在配置文件celeryconfig.py中我们已经开启了处理的结果回应模式了celery_ignore_result = false并且在回应方式配置中我们设置了celery_result_backend = 'amqp'这样我们就可以查看到处理的状态了。

  1.  
    >>> x = app.send_task('task.say',['hello', 'lady'])
  2.  
    >>> x.ready()
  3.  
    false
  4.  
    >>> x.status
  5.  
    'pending'
  6.  
    >>> x.ready()
  7.  
    true
  8.  
    >>> x.status
  9.  
    u'success'

可以看到任务发送给celery后马上查看任务状态会处于pending状态。稍等片刻就可以查看到success状态了。这种效果真棒不是吗?在图像处理中或者其他的一些搞耗时的任务中,我们只需要把任务发送给后台就不用去管它了。当我们需要结果的时候只需要查看一些是否成功完成了,如果返回成功我们就可以去后台数据库去找处理后生成的数据了。

 

celery使用mangodb保存数据

第一步

安装好mongodb了!就可以使用它了,首先让我们修改celeryconfig.py文件,使celery知道我们有一个新成员要加入我们的项目,它就是mongodb配置的方式如下。

elery_imports = ('tasks')
celery_ignore_result = false
broker_host = '127.0.0.1'
broker_port = 5672
broker_url = 'amqp://'
#celery_result_backend = 'amqp'
celery_result_backend = 'mongodb'
celery_result_backend_settings = {
        "host":"127.0.0.1",
        "port":27017,
        "database":"jobs",
        "taskmeta_collection":"stock_taskmeta_collection",
}

把#celery_result_backend = 'amp'注释掉了,但是没有删除目的是对比前后的改变。为了使用mongodb我们有简单了配置一下主机端口以及数据库名字等。显然你可以按照你喜欢的名字来配置它。

第二步

启动 mongodb 数据库:mongod。修改客户端client.py让他能够动态的传人我们的数据,非常简单代码如下。

  1.  
    import sys
  2.  
    from celery import celery
  3.  
     
  4.  
    app = celery()
  5.  
     
  6.  
    app.config_from_object('celeryconfig')
  7.  
    app.send_task("tasks.say",[sys.argv[1],sys.argv[2]])

任务tasks.py不需要修改!

  1.  
    import time
  2.  
    from celery.task import task
  3.  
     
  4.  
     
  5.  
    @task
  6.  
    def say(x,y):
  7.  
    time.sleep(5)
  8.  
    return x+y

第三步

测试代码,先启动celery任务。

celery worker -l info --beat

再来启动我们的客户端,注意这次启动的时候需要给两个参数啦!
mongo

python client.py welcome landpack

等上5秒钟,我们的后台处理完成后我们就可以去查看数据库了。

第四步

查看mongodb,需要启动一个mongodb客户端,启动非常简单直接输入 mongo 。然后是输入一些简单的mongo查询语句。

最后查到的数据结果可能是你不想看到的,因为mongo已经进行了处理。想了解更多可以查看官方的文档。

上一篇:

下一篇: