【发布时间】:2013-08-31 09:09:27
【问题描述】:
我正在学习 AMQP,所以这可能是对我在做什么的误解。我正在尝试在我自己的私人服务器上设置rabbitmq,以便为一组系统提供服务(我有一堆需要处理的图像)。我设置了一个三步管道:
| put image on queue | --work_queue--> | process | --results--> | archive results |
我安装了rabbitmq,并在服务器上创建了两个队列。 “工作队列”和“结果”。我使用 amqplib 编写了一些 python 脚本,并且我的管道与一个图像处理器工作人员一起工作得很好。我已将 100 张图像添加到队列中,而我的一台机器愉快地一次抓取一张并处理数据,并将结果放入结果队列中。
问题是,我假设如果我在另一台机器上启动另一个图像处理程序,它也会简单地将工作从队列中拉出。这似乎是 rabbitmq 网站上的“工作队列”教程中列出的确切情况。我希望这能正常工作。然而,真正发生的是,无论我启动多少其他工作人员,他们只是永远等待并且永远不会得到任何工作,即使在 work_queue 中有大量消息在等待。
我误解了什么?在服务器上排队工作的相关代码:
from amqplib import client_0_8 as amqp
conn = amqp.Connection(host="foo:5672 ", userid="pipeline",
password="XXXXX", virtual_host="/", insist=False)
chan = conn.channel()
....
msg = { 'filename': os.path.basename(i) }
chan.basic_publish(amqp.Message(json.dumps(msg)), exchange='',
routing_key='work_queue')
以及流程工作者的消费者方面:
from amqplib import client_0_8 as amqp
conn = amqp.Connection(host="foo:5672 ", userid="pipeline",
password="XXXXX", virtual_host="/", insist=False)
chan = conn.channel()
def work_callback(msg):
....
while True:
chan.basic_consume(callback=work_callback, queue='work_queue')
try:
chan.wait()
except KeyboardInterrupt:
print "\nExiting cleanly"
sys.exit(0)
我可以看到其他工作人员已连接:
$ sudo rabbitmqctl list_queues
Listing queues ...
results 0
work_queue 246
...done.
$ sudo rabbitmqctl list_connections
Listing connections ...
pipeline 192.168.8.1 41553 running
pipeline XX.YY.ZZ.WW 46676 running
pipeline 192.168.8.4 44482 running
pipeline 192.168.8.6 41884 running
...done.
其中 XX.YY.ZZ.WW 是外部 IP。位于 192.168.8.6 的工作人员正在排队等待,但外部 IP 上的工作人员闲置在那里,而 246 条消息在它应该等待的队列中等待。
想法?
【问题讨论】:
-
令人着迷。如果在我提交工作时两个“工作人员”都已连接,那么两个工作人员都可以工作。甚至更奇怪的是,如果在提交工作时只有一个连接,就好像所有工作都安排给那个工人..如果另一个工人连接,它只有在我杀死原来的工人时才开始工作......然后是新的一个人得到了大约一半的工作???如果我重新连接原来的,它会得到剩下的一半?奇怪!
-
我希望工作被保存在一个队列中,并且每个工作人员连接到队列并一次获取第一个条目。相反,它似乎将所有工作安排给它。也许我不明白队列和渠道?