【问题标题】:How to safely close a pika BlockingConnection subscriber?如何安全地关闭 pika BlockingConnection 订阅者?
【发布时间】:2021-01-10 13:43:48
【问题描述】:
根据文档Pika is not thread safe:
Pika 在代码中没有任何线程概念。如果您想将 Pika 与线程一起使用,请确保每个线程都有一个在该线程中创建的 Pika 连接。跨线程共享一个 Pika 连接是不安全的,但有一个例外:您可以从另一个线程调用连接方法 add_callback_threadsafe 以在活动的 pika 连接中安排回调。
假设我有一个订阅者,我已经开始使用channel.start_consuming()。该线程将被阻止等待消息到达。这些消息可能相隔很长时间(有时是几个小时)。
当然,如果我想安全/干净地关闭订阅者,我必须从另一个线程这样做吗?否则如何触发消费者突破阻塞?
【问题讨论】:
标签:
multithreading
amqp
pika
【解决方案1】:
您可以使用connection.process_data_events() 而不仅仅是channel.start_consuming()。这里的好处是你可以做这样的事情来关闭连接。
def consume_messages(self):
while self.running:
self.connection.process_data_events()
sleep(0.1)
self.connection.close()
然后您只需将self.running 设置为False 即可关闭连接。