【问题标题】:How can I run both a consumer and a publisher at the same time in Pika?如何在 Pika 中同时运行消费者和发布者?
【发布时间】:2016-01-19 21:36:18
【问题描述】:

我有一个 Python 应用程序,我想同时运行消费者和发布者。基本上,我想通过Consumer获取消息,然后对其进行解析和处理,然后通过Publisher将其发送回RabbitMQ。

我从官方 Pika 文档中获取了 async consumerasync publisher 的代码。它们单独工作,但我似乎无法让它们同时工作。在应用程序的起点,我有:

import messaging.MQ as MQ

MQ.start_consumer()
MQ.start_publisher()

但是,永远无法到达start_publisher() 行。

看起来罪魁祸首是Consumer中的这一行:

self._connection.ioloop.start()

在调用该行之后,除了Consumer 中定义的 asnyc 方法之外,什么都不会执行。

我觉得这很明显,但我似乎无法确定。

【问题讨论】:

  • 你必须将发布逻辑放在你自己的on_message回调中
  • @0x41ndrea:这可能可行,但需要进行大量重构。请在下面查看我的答案。

标签: python rabbitmq pika


【解决方案1】:

这有点尴尬,但我在发布后 10 分钟内解决了这个问题。

基本上,您只需要使用线程即可。可能有更简单的方法,但这确实有效。

按照文档中的示例,以Consumer 为例,您只需替换:

class ExampleConsumer(object):

与:

import threading

class ExampleConsumer(threading.Thread):

然后你可以简单地运行线程:

publisher = Publisher()
consumer = Consumer()

publisher.start()
consumer.start()

我还调整了 __init__ 函数并添加了一些 Exception 捕捉,但试图在代码示例中保持简单。

【讨论】:

  • 确实,您的类现在继承了threading.thread 类,其start() 方法打开了一个新线程并执行run() 方法。您之前实现的问题是您在当前线程中运行消费者。消费者进入一个无限循环(等待消息事件),所以当然不会执行超过该行的任何内容。
  • 是的,事后看来,这是非常明显的。我只是假设 Pika 实现会以某种方式克服这一点,因为我猜很少有人真正希望他们的 RabbitMQ 消费者阻止其他一切。无论如何,这很容易过去,希望它将来也可以帮助其他人。谢谢@ktbiz!
  • 在多线程模式中start()打开一个新线程并执行run(),这是用户实现的工作函数(通常是阻塞的)。您不应该发现自己实现了start() 方法。您实现run(),然后在包含的线程对象上调用start()。在这种情况下,您已将线程子类化,但您也可以实例化一个 thread 对象,将实现 run() 的对象传递给它,然后在线程上调用 start(),这将分叉然后在您的线程上调用 run()目的。例如:t = threading.Thread(myConsumer); t.start()
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-06-10
相关资源
最近更新 更多