【问题标题】:How to detect a queue has been deleted?如何检测队列已被删除?
【发布时间】:2017-03-01 20:43:30
【问题描述】:

当我手动删除我的 PikaClient 使用的队列时,什么也没有发生。我可以重新创建具有相同名称的队列,但通道已停止使用队列(正常,因为我已将其删除)。但是我想在消费队列被删除时收到一个事件。

我预计频道会自动关闭,但从未调用过 «on_channel_close_callback»。 «basic_consume» 在关闭时不提供任何回调。 还有一点很重要,我必须使用 TornadoConnection。

鼠兔:0.10.0 蟒蛇:2.7 龙卷风:4.3

谢谢你的帮助。

class PikaClient(object):

    def __init__(self):
        # init everything here

    def connect(self):
        pika.adapters.tornado_connection.TornadoConnection(connection_param, on_open_callback=self.on_connected)

    def on_connected(self, connection):
        self.logger.info('PikaClient: connected to RabbitMQ')
        self.connected = True
        self.connection = connection
        self.connection.channel(self.on_channel_open)

    def on_open_error_callback(self, *args):
        self.logger.error("on_open_error_callback")

    def on_channel_open(self, channel):
        channel.add_on_close_callback(self.on_channel_close_callback)

        channel.basic_consume(self.on_message, queue=self.queue_name, no_ack=True)

    def on_channel_close_callback(self, reply_code, reply_text):
        self.logger.error("Consumer was cancelled remotely, shutting down", reply_code=reply_code, reply_text=reply_text)

【问题讨论】:

    标签: python rabbitmq tornado pika


    【解决方案1】:

    我找到了解决方法。 如果我的 PikaClient 已使用消息,我会每 X 秒检查一次。如果没有,我重新启动将自动创建队列的应用程序。

    如果您有更好的解决方案,我仍然愿意提供建议。

    def __init__(self):
        ...
        self.have_messages_been_consumed = False
    
    def on_connected(self, connection):
        self.logger.info('PikaClient: connected to RabbitMQ')
        self.connected = True
        self.connection = connection
        self.connection.add_timeout(X, self.check_if_messages_have_been_consumed)
        self.connection.channel(self.on_channel_open)
    
    def check_if_messages_have_been_consumed(self):
        if self.have_messages_been_consumed:
            self.have_messages_been_consumed = False
            self.connection.add_timeout(X, self.check_if_messages_have_been_consumed)
        else:
            # close_and_restart will set to False have_messages_been_consumed
            self.close_and_restart()
    
    def on_message(self, channel, basic_deliver, header, body):
        self.have_messages_been_consumed = True
        ...
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-08-24
      • 2012-03-10
      • 1970-01-01
      • 1970-01-01
      • 2017-01-27
      • 1970-01-01
      • 1970-01-01
      • 2018-04-18
      相关资源
      最近更新 更多