You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RabbitMQ未ACK消费大量队列消息异常问题及解决方案咨询

问题描述

需求:从RabbitMQ队列读取100,000条消息后执行Negative Acknowledge(NACK),让其他消费者后续ACK并处理。使用以下Python代码尝试打印所有消息信息:

# Script settings
seen_message_ids = set()
retrieved_number = 0
delete_all_messages = False
ids_to_delete = ['da140a51-cdd9-f71d-ac5c-ec15766e6aa4']

def on_message(ch, method_frame, header_frame, body):
    # If all arguments other than `ch` are None, then the `consume` method timed out due to queue inactivity.
    # Stop listening to the queue.
    if method_frame is None and header_frame is None and body is None:
        ch.stop_consuming()
        return

    # Messages that are not deleted are re-queued, meaning they will be consumed again, infinitely. Keep track
    # of `message_id`, which is a UUID to identify message that were already seen. Since this is a FIFO queue,
    # seeing any message again means all messages in the queue have been examined.
    message_id = header_frame.message_id
    if message_id in seen_message_ids:
        print("All messages evaluated.")
        ch.stop_consuming()
        return
    seen_message_ids.add(message_id)

    # Keep track of the number of retrieved messages.
    global retrieved_number
    retrieved_number += 1

    # To keep the message, negatively acknowledge it and re-queue it.
    ch.basic_nack(delivery_tag=method_frame.delivery_tag, requeue=True)
    status = "READ   "
    print(f"{status} {retrieved_number:06d} Message: {body.decode()} UUID: {header_frame.message_id}")
    return


def main():
    connection = pika.BlockingConnection(
            pika.ConnectionParameters(host=RABBITMQ_HOST,
                                      port=RABBITMQ_PORT,
                                      virtual_host=RABBITMQ_VHOST,
                                      credentials=credentials))
    channel = connection.channel()
    channel.queue_declare(queue='test1')

    # Consume messages and acknowledge them
    print("START LISTENING")
    last_message_received_time = datetime.datetime.now()
    for method_frame, properties, body in channel.consume('test1', inactivity_timeout=1):

        # `consume` is a generator https://pika.readthedocs.io/en/stable/examples/blocking_consumer_generator.html
        # That means that the body of this for loop will not be reached unless a new message is found.
        new_message_received_time = datetime.datetime.now()
        if method_frame is None and properties is None and body is None:
            print(f"Timeout due to inactivity after {new_message_received_time - last_message_received_time} elapsed.")
        on_message(channel, method_frame, properties, body)

    try:
        channel.start_consuming()
    except KeyboardInterrupt:
        channel.stop_consuming()
    finally:
        print("STOP LISTENING")
        connection.close()

执行后出现以下异常:

  • 打印约5000条消息后输出"All messages evaluated",但所有message_id都是唯一UUID,说明重入队的消息被提前消费,不符合FIFO队列预期(重入队消息应排在原有10万条之后)
  • RabbitMQ管理UI显示大部分消息处于Unacked状态,需数分钟才能归零
  • 消费者卡在ch.stop_consuming()处,直到所有消息变为Ready
  • 第二个消费者或管理UI无法获取消息

问题原因分析
  1. 预取机制导致消息批量堆积:RabbitMQ默认prefetch_count为0,Broker会尽可能多地将队列消息推送给消费者(直到TCP缓冲区耗尽)。消费者一次性获取约5000条消息并标记为Unacked,随后逐个NACK重入队。此时Ready队列中已存在这些重入队消息,消费者处理完手中批次后,会优先读取这些重复消息,触发代码中“message_id重复”的判断,提前终止消费,遗漏大量未处理的Unacked消息。

  2. Unacked消息延迟释放:调用stop_consuming后,RabbitMQ需要将消费者手中所有未确认的Unacked消息重新标记为Ready,这个过程因消息数量大需要耗时,导致消费者卡在该步骤直到所有消息回到Ready状态。

  3. 队列资源被独占:大量消息处于Unacked状态时,其他消费者或管理UI无法获取这些消息,因为它们被当前消费者持有,只有回到Ready队列后才能被再次消费。


解决办法

1. 限制预取数量,避免批量堆积

设置合理的prefetch_count,让消费者每次仅获取少量消息(如100条),确保处理完当前批次后再获取新消息,避免重入队消息被提前消费:

def main():
    # 省略连接代码
    channel = connection.channel()
    channel.queue_declare(queue='test1')
    # 设置预取数量,控制单次获取的消息数
    channel.basic_qos(prefetch_count=100)
    # 后续消费逻辑不变

2. 调整终止判断逻辑,基于消息总数控制消费

取消基于message_id重复的终止逻辑,改为通过查询队列总消息数来判断是否完成所有消息的读取:

def main():
    connection = pika.BlockingConnection(...)
    channel = connection.channel()
    # 被动查询队列信息,获取总消息数
    queue_result = channel.queue_declare(queue='test1', passive=True)
    total_messages = queue_result.method.message_count
    channel.basic_qos(prefetch_count=100)

    retrieved_number = 0
    print(f"START LISTENING, total messages: {total_messages}")

    for method_frame, properties, body in channel.consume('test1', inactivity_timeout=5):
        if method_frame is None:
            print("Timeout, no more messages to consume")
            break
        
        retrieved_number += 1
        # NACK并重入队
        channel.basic_nack(delivery_tag=method_frame.delivery_tag, requeue=True)
        print(f"READ   {retrieved_number:06d} Message: {body.decode()} UUID: {properties.message_id}")
        
        # 达到总消息数后终止消费
        if retrieved_number >= total_messages:
            print("All messages evaluated.")
            break

    # 关闭连接
    channel.cancel()
    connection.close()
    print("STOP LISTENING")

注意:passive=True的queue_declare仅查询队列信息,不会创建队列,适合获取现有队列的消息数。

3. 优化消息流转逻辑,避免重复循环

如果无需立即将消息放回原队列,可将消息转发到临时队列暂存,再ACK原消息,避免原队列出现消息循环:

def on_message(ch, method_frame, header_frame, body):
    # 省略其他逻辑
    # 将消息转发到临时队列
    ch.basic_publish(exchange='', routing_key='temp_queue', body=body, properties=header_frame)
    # ACK原消息,从原队列移除
    ch.basic_ack(delivery_tag=method_frame.delivery_tag)

后续消费者可从temp_queue中获取消息进行处理。

4. 简化消费逻辑,避免冲突

原代码同时使用了channel.consume生成器和channel.start_consuming(),存在逻辑冲突,建议只保留一种消费方式(如仅使用生成器循环),移除多余的start_consuming调用。


内容的提问来源于stack exchange,提问作者Jacob Quisenberry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 15:17:33