代码之家  ›  专栏  ›  技术社区  ›  verisimilitude

RabbitMQ未确认消息未得到重新请求

  •  0
  • verisimilitude  · 技术社区  · 8 年前

    我在长期运行RabbitMQ消费者时面临一个问题。我的几条消息最终都处于未确认状态。

    我的RabbitMQ版本:3.6.15 Pika版本:0.11.0b

    import pika
    import time
    import sys
    import threading
    from Queue import Queue
    rabbitmq_server = "<SERVER>"
    queue = "<QUEUE>"
    connection = None
    
    def check_acknowledge(channel, connection, ack_queue):
        delivery_tag = None
        while(True):
            try:
                delivery_tag = ack_queue.get_nowait()
                channel.basic_nack(delivery_tag=delivery_tag)
                break
            except:
                connection.process_data_events()
            time.sleep(1)
    
    
    def process_message(body, delivery_tag, ack_queue):
        print "Received %s" % (body)
        print "Waiting for 600 seconds before receiving next ID\n"
        start = time.time()
        elapsed = 0
        while elapsed < 10:
            elapsed = time.time() - start
            print "loop cycle time: %f, seconds count: %02d" %(time.clock(), elapsed)
            time.sleep(1)
        ack_queue.put(delivery_tag)
    
    
    
    
    def callback(ch, method, properties, body):
        global connection
        ack_queue = Queue()
        t = threading.Thread(target=process_message, args=(body, method.delivery_tag, ack_queue))
        t.start()
        check_acknowledge(ch, connection, ack_queue)
    
    while True:
        try:
            connection = pika.BlockingConnection(pika.ConnectionParameters(host=rabbitmq_server))
            channel = connection.channel()
            print ' [*] Waiting for messages. To exit press CTRL+C'
            channel.basic_qos(prefetch_count=1)
            channel.basic_consume(callback, queue=queue)
            channel.start_consuming()
        except KeyboardInterrupt:
            break
    
    channel.close()
    connection.close()
    exit(0)
    

    我错过什么了吗?

    1 回复  |  直到 8 年前
        1
  •  0
  •   verisimilitude    6 年前

    我使用以下多线程使用者来解决此问题。

    import pika
    import time
    import sys
    import threading
    from Queue import Queue
    rabbitmq_server = "<RABBITMQ_SERVER_IP>"
    queue = "hello1"
    connection = None
    
    
    
    
    def check_acknowledge(channel, connection, ack_queue):
        delivery_tag = None
        while(True):
            try:
                delivery_tag = ack_queue.get_nowait()
                channel.basic_ack(delivery_tag=delivery_tag)
                break
            except:
                connection.process_data_events()
            time.sleep(1)
    
    
    def process_message(body, delivery_tag, ack_queue):
        print "Received %s" % (body)
        print "Waiting for 600 seconds before receiving next ID\n"
        start = time.time()
        elapsed = 0
        while elapsed < 300:
            elapsed = time.time() - start
            print "loop cycle time: %f, seconds count: %02d" %(time.clock(), elapsed)
            time.sleep(1)
        ack_queue.put(delivery_tag)
    
    
    
    
    def callback(ch, method, properties, body):
        global connection
        ack_queue = Queue()
        t = threading.Thread(target=process_message, args=(body, method.delivery_tag, ack_queue))
        t.start()
        check_acknowledge(ch, connection, ack_queue)
    
    while True:
        try:
            connection = pika.BlockingConnection(pika.ConnectionParameters(host=rabbitmq_server))
            channel = connection.channel()
            print ' [*] Waiting for messages. To exit press CTRL+C'
            channel.basic_qos(prefetch_count=1)
            channel.basic_consume(callback, queue=queue)
            channel.start_consuming()
        except KeyboardInterrupt:
            break
    
    channel.close()
    connection.close()
    exit(0)
    
    1. 消费者 callback 函数触发单独的函数 check_acknowledge 在主线程本身中。因此,连接和通道对象保留在同一线程中。请注意,Pika不是线程安全的,因此我们需要在同一线程中维护这些对象。
    2. 实际处理发生在从主线程派生的新线程中。
    3. 一旦 process_message 处理完成后 delivery_tag 在队列中。

    4. check\u确认 无限循环,直到找到 delivery\u标签 排队人 process\u消息 。一旦找到,它就会 acks 消息并返回。

    我已通过运行此使用者测试了此实现 sleep 持续5分钟、10分钟、30分钟和1小时。这对我来说很有效。

    推荐文章