【发布时间】:2018-07-18 15:02:20
【问题描述】:
我有一个用例,我需要启动一个 celery 工作人员,以便他们使用唯一的队列,我尝试如下实现。
from celery import Celery
app = Celery(broker='redis://localhost:9555/0')
@app.task
def afunction(arg1=None, arg2=None, arg3=None):
if arg1 == 'awesome_1':
return "First type of Queue executed"
if arg2 == "awesome_2":
return "Second Type of Queue executed"
if arg3 == "awesome_3":
return "Third Type of Queue executed"
if __name__=='__main__':
qlist = ["awesome_1", "awesome_2", "awesome_3"]
arglist = [None, None, None]
for q in qlist:
arglist[qlist.index(q)] = q
argv = [
'worker',
'--detach',
'--queue={0}'.format(q),
'--concurrency=1',
'-E',
'--loglevel=INFO'
]
app.worker_main(argv)
afunction.apply_async(args=[arglist[0], arglist[1], arglist[2]], queue=q)
此代码在执行时给出以下输出:
[2018-02-08 11:28:43,479: INFO/MainProcess] Connected to redis://localhost:9555/0
[2018-02-08 11:28:43,486: INFO/MainProcess] mingle: searching for neighbors
[2018-02-08 11:28:44,503: INFO/MainProcess] mingle: all alone
[2018-02-08 11:28:44,527: INFO/MainProcess] celery@SYSTEM ready.
[2018-02-08 11:28:44,612: INFO/MainProcess] Received task: __main__.afunction[f092f721-6523-4055-98fc-158ac316f4cc]
[2018-02-08 11:28:44,618: INFO/ForkPoolWorker-1] Task __main__.afunction[f092f721-6523-4055-98fc-158ac316f4cc] succeeded in 0.0010992150055244565s: 'First type of Queue executed'
因此,我可以看到工作程序在 for 循环的第一次迭代中按应有的方式执行,但随后它就停在那里并且不再继续执行 for 循环。
我相信这是因为 worker 没有独立运行,或者作为脚本的子进程运行,因为我可以看到 1 + 与 python 在ps aux 上运行相同脚本的进程一样多,因为正在设置--concurrency。关于出了什么问题或如何使工作队列分离运行的任何指针,因此在return 之后for 循环继续迭代。
【问题讨论】:
标签: python multithreading celery message-queue