架构:
这里作为例子的celery app为myapi:
公共代码部分:
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)
发送任务给machine1:
用get()可以看到来自machine1的返回,再看看machine1的显示:
另外的machine2的节点部署和machine1的一样,启动workers的时候指定的queue为machine2
在发布任务的时候指定queue和routing_key到不同节点上,就可以将任务分配到不同的节点上运行