【问题标题】:How to restart consumer rabbitmq pika python如何重启消费者rabbitmq pika python
【发布时间】:2016-05-12 09:36:22
【问题描述】:

我在使用 rabbitmq 的 pika 库设置的消费者中遇到了一些丢失。与 pika 一起,我使用扭曲的实现来设置异步消费者。我不确定为什么会发生这种情况,但如果消费者退出并且不确定如何执行此操作,我希望实现重新连接。这是我目前的实现

class Consumer(object):
def __init__(self, queue, exchange, routingKey, medium, signalRcallbackFunc):
    self._queue_name = queue
    self.exchange = exchange
    self.routingKey = routingKey
    self.medium = medium
    print "client on"
    self.channel = None
    self.medium.client.on(signalRcallbackFunc, self.callback)

def on_connected(self, connection):
    d = connection.channel()
    d.addCallback(self.got_channel)
    d.addCallback(self.queue_declared)
    d.addCallback(self.queue_bound)
    d.addCallback(self.handle_deliveries)
    d.addErrback(log.err)

def got_channel(self, channel):
    self.channel = channel
    self.channel.basic_qos(prefetch_count=500)
    return self.channel.queue_declare(queue=self._queue_name, durable=True)

def queue_declared(self, queue):
    self.channel.queue_bind(queue=self._queue_name,
                            exchange=self.exchange,
                            routing_key=self.routingKey)

def queue_bound(self, ignored):
    return self.channel.basic_consume(queue=self._queue_name)

def handle_deliveries(self, queue_and_consumer_tag):
    queue, consumer_tag = queue_and_consumer_tag
    self.looping_call = task.LoopingCall(self.consume_from_queue, queue)

    return self.looping_call.start(0)

def consume_from_queue(self, queue):
    d = queue.get()
    return d.addCallback(lambda result: self.handle_payload(*result))

def handle_payload(self, channel, method, properties, body):
    print(body)
    print(properties.headers)
    channel.basic_ack(method.delivery_tag)
    print "#####################################" + method.delivery_tag + "###################################"

def callback(self, data):
    #self.channel.basic_ack(data, multiple=True)
    pass

【问题讨论】:

    标签: python python-2.7 rabbitmq twisted pika


    【解决方案1】:

    您可以在 on_connected 回调中为连接注册一个“关闭”处理程序。这在连接丢失时被调用。在这里,您可以重新建立新的连接。

    下面的例子比较有用,是我用过的一个策略,效果不错…… http://pika.readthedocs.io/en/latest/examples/asynchronous_consumer_example.html

    对于twisted pika 库,add_on_close_callback 方法可能会让你走得很远(尽管我还没有测试过)。 https://pika.readthedocs.io/en/0.10.0/modules/adapters/twisted.html

    【讨论】:

    • 我将如何使用扭曲协议进行此操作
    • 看起来 pika 库的扭曲实现有类似的方法可用于注册连接关闭事件的侦听器。您可以尝试使用add_on_close_callback 方法。见[链接]pika.readthedocs.io/en/0.10.0/modules/adapters/twisted.html
    • 我已阅读此内容,但不确定如何在我的解决方案中实现此内容
    【解决方案2】:

    你有什么理由不能关闭连接并重新打开它?

    @contextmanager
    def with_pika_connection():
        credentials = pika.PlainCredentials(worker_config.username, worker_config.password)
        connection = pika.BlockingConnection(pika.ConnectionParameters(
            host=worker_config.host,
            credentials=credentials,
            port=worker_config.port,
        ))
    
        try:
            yield connection
        finally:
            connection.close()
    
    
    @contextmanager
    def with_pika_channel(queuename):
        with with_pika_connection() as connection:
            channel = connection.channel()
    
    
    while True:
        while not stopping:
             try:
                    with with_pika_channel(queuename) as (connection, channel):
                        consumer_tag = channel.basic_consume(
                            callback,
                            queue=queuename,
                        )
                        channel.start_consuming()
             except Exception as e:
                  reportException(e) 
                  # Continue 
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-10-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多