celery集群

celery devops 阅读:1377

架构:

blob.png

这里作为例子的celery app为myapi:


blob.png


公共代码部分:

celery.py:

from __future__ import absolute_import
from celery import Celery
app = Celery('myapi',
             broker='redis://127.0.0.1:32189/0',
             backend='redis://127.0.0.1:32189/1',
             include=['myapi.agent'])
app.config_from_object('myapi.config')
if __name__ == '__main__':
  app.start()

config.py:

from __future__ import absolute_import
from kombu import Queue,Exchange
from datetime import timedelta
CELERY_TASK_RESULT_EXPIRES=3600
CELERY_TASK_SERIALIZER='json'
CELERY_ACCEPT_CONTENT=['json']
CELERY_RESULT_SERIALIZER='json'
CELERYD_CONCURRENCY = 5
CELERY_DEFAULT_EXCHANGE = 'agent'
CELERY_DEFAULT_EXCHANGE_TYPE = 'direct'
CELERT_QUEUES =  (
  Queue('machine1',exchange='agent',routing_key='machine1'),
  Queue('machine2',exchange='agent',routing_key='machine2'),
)

__init__.py:(空白)


任务节点agent.py:

from __future__ import absolute_import
from myapi.celery import app
@app.task
def add(x,y):
    return {'the value is ':str(x+y)}
@app.task
def writefile():
    out=open('/tmp/data.txt','w')
    out.write('hello'+'\n')
    out.close()
@app.task
def mul(x,y):
    return x*y
@app.task
def xsum(numbers):
    return sum(numbers)
@app.task
def getl(stri):
    return getlength(stri)
def getlength(stri):
    return len(stri)


在这个例子中只测试mul()函数:

在myapi几点上启动worker:(用-Q指定监听的queue)


blob.png


发送任务给machine1:

blob.png


用get()可以看到来自machine1的返回,再看看machine1的显示:

blob.png


另外的machine2的节点部署和machine1的一样,启动workers的时候指定的queue为machine2

在发布任务的时候指定queue和routing_key到不同节点上,就可以将任务分配到不同的节点上运行