【问题标题】:AMQP / RabbitMQ multiple consumers but only one geting work?AMQP / RabbitMQ 多个消费者但只有一个工作?
【发布时间】: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 条消息在它应该等待的队列中等待。

想法?

【问题讨论】:

  • 令人着迷。如果在我提交工作时两个“工作人员”都已连接,那么两个工作人员都可以工作。甚至更奇怪的是,如果在提交工作时只有一个连接,就好像所有工作都安排给那个工人..如果另一个工人连接,它只有在我杀死原来的工人时才开始工作......然后是新的一个人得到了大约一半的工作???如果我重新连接原来的,它会得到剩下的一半?奇怪!
  • 我希望工作被保存在一个队列中,并且每个工作人员连接到队列并一次获取第一个条目。相反,它似乎将所有工作安排给它。也许我不明白队列和渠道?

标签: python rabbitmq amqp


【解决方案1】:

很抱歉回答我自己的问题,但我终于得到了我正在寻找的行为。我需要将 QoS prefetch_count 设置为 1。显然,通过不设置 prefetch_count,只要一个客户端连接,所有 200 多条消息都会被传递到它的“通道”并从队列中排出,所以除非其他人可以看到任何工作,否则原来的工作断开连接(因此关闭通道并将消息放回队列。

通过添加:

chan.basic_qos(0,1,False)

对于我的员工来说,他们现在一次只能抓取一条消息。不完全是我所期望的,但它现在仍然按预期工作。

【讨论】:

    猜你喜欢
    • 2020-09-12
    • 2012-05-24
    • 2017-04-13
    • 1970-01-01
    • 2015-12-28
    • 2023-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多