【问题标题】:Consume multiple queues in python / pika在 python / pika 中消费多个队列
【发布时间】:2014-08-22 00:23:09
【问题描述】:

我正在尝试创建一个订阅多个队列的消费者,然后在消息到达时对其进行处理。

问题在于,当第一个队列中已经存在一些数据时,它会消耗第一个队列,而永远不会去消耗第二个队列。 但是,当第一个队列为空时,它确实会转到下一个队列,然后同时消耗两个队列。

我首先实现了线程,但想避开它,当 pika 库为我完成它时没有太多复杂性。以下是我的代码:

import pika

mq_connection = pika.BlockingConnection(pika.ConnectionParameters('x.x.x.x'))
mq_channel = mq_connection.channel()
mq_channel.basic_qos(prefetch_count=1)


def callback(ch, method, properties, body):
    print body
    mq_channel.basic_ack(delivery_tag=method.delivery_tag)

mq_channel.basic_consume(callback, queue='queue1', consumer_tag="ctag1.0")
mq_channel.basic_consume(callback, queue='queue2', consumer_tag="ctag2.0")
mq_channel.start_consuming()

【问题讨论】:

  • 我尝试了您的代码,唯一的更改是添加记录器以防止异常,并声明队列。该代码按预期工作。我向每个队列发布了一些消息,这些消息在 CLI 上被路由和回显
  • 您好,您可以尝试使用预先填充的队列,然后启动消费者。让我知道这是否也能按预期工作。
  • 我刚刚试过了,它不起作用。我只看到来自第一个队列的消息。
  • 这就是我要说的。是不是很奇怪?你有什么想法吗?
  • 我对python客户端不太了解,所以才请Gavin bellow回答

标签: python rabbitmq pika


【解决方案1】:

值得注意的是,上述解决方案仅在 auto_ack 设置为 True 时才有效。

【讨论】:

    【解决方案2】:

    与上面第一个答案中的 cmets 类似,我能够使用 pika 1.1.0 和以下版本获得类似的结果:

    import pika
    
    def queue1_callback(ch, method, properties, body):
      print(" [x] Received queue 1: %r" % body)
    
    def queue2_callback(ch, method, properties, body):
      print(" [x] Received queue 2: %r" % body)
    
    def on_open(connection):
      connection.channel(on_open_callback = on_channel_open)
    
    
    def on_channel_open(channel):
      channel.basic_consume('queue1', queue1_callback, auto_ack = True)
      channel.basic_consume('queue2', queue2_callback, auto_ack = True)
    
    credentials = pika.PlainCredentials('u', 'p')
    parameters = pika.ConnectionParameters('localhost', 5672, '/', credentials)
    connection = pika.SelectConnection(parameters = parameters, on_open_callback = on_open)
    
    Try:
      connection.ioloop.start()
    except KeyboardInterrupt:
      connection.close()
      connection.ioloop.start()

    【讨论】:

      【解决方案3】:

      一种可能的解决方案是使用非阻塞连接并使用消息。

      import pika
      
      
      def callback(channel, method, properties, body):
          print(body)
          channel.basic_ack(delivery_tag=method.delivery_tag)
      
      
      def on_open(connection):
          connection.channel(on_channel_open)
      
      
      def on_channel_open(channel):
          channel.basic_consume(callback, queue='queue1')
          channel.basic_consume(callback, queue='queue2')
      
      
      parameters = pika.URLParameters('amqp://guest:guest@localhost:5672/%2F')
      connection = pika.SelectConnection(parameters=parameters,
                                         on_open_callback=on_open)
      
      try:
          connection.ioloop.start()
      except KeyboardInterrupt:
          connection.close()
      

      这将连接到多个队列并相应地使用消息。

      【讨论】:

      • 你能告诉我最后 %2F 的用途吗?
      • @RápliAndrás 连接rabbitmq时,需要指定virtualhost。默认主机为/,转义为%2f
      • 值得注意的是,这段代码不适用于 Pika 1.1.0。只需在 on_open 方法中添加 on_open_callback=:connection.channel(on_open_callback=on_channel_open) 和 on_channel_open 方法中的 on_message_callback=:channel.basic_consume(on_open_callback=callback, queue='queue1') channel.basic_consume(on_open_callback=callback , queue='queue2')
      • 有没有办法定义队列之间的优先级?
      【解决方案4】:

      这个问题很可能是第一个调用发出了 Basic.Consume 并且在发出第二个调用之前已经从预先填充的队列中接收到消息。您可能想尝试将 QoS 预取计数设置为 1,这将限制 RabbitMQ 一次向您发送多条消息。

      【讨论】:

      • 它已经被设置为1,从代码中可以看出。你还能想到什么?
      • 还有一件事,我认为消费者在达到 start_sumption 之前不会真正开始消费。不过需要验证。
      • 嗨,Gavin,我刚刚尝试使用 python 的 kombu 库来实现这个功能,它按预期工作。
      • 我会尝试使用 Async 后端并通知您。
      • 我关注了pika.readthedocs.org/en/0.9.13/examples/…,一切正常。第一个队列被完全消耗,然后第二个队列被相应地消耗。
      猜你喜欢
      • 2023-03-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多