Asyncio对接RabbitMQ:如何判断是否还有待消费的RabbitMQ消息
解决方案
你使用的queue.consume()属于推模式消费,aio-pika会先把RabbitMQ推送过来的消息存在客户端内部缓冲队列再调度回调,你可以通过以下两种方式实现需求:
方案1:计数器+QoS预取限制(推荐,稳定兼容)
这是符合RabbitMQ最佳实践的实现方式,无私有属性依赖:
- 首先给channel设置预取数,和你预期的单批次消息数量一致,控制RabbitMQ一次只会推送指定数量的消息到客户端,不会多推
- 维护两个变量统计处理状态:
pending_msg_count代表已经接收待处理的消息数量,处理完成后减1 - 每次消息处理完成后检查
pending_msg_count是否为0,为0就说明当前所有已接收的消息都处理完成,可以触发通知逻辑
代码示例:
import aio_pika # 也可以用类属性封装,避免全局变量 pending_msg_count = 0 # 你单批次的消息处理数量,可自行调整 BATCH_SIZE = 10 async def subscribe_to_messages(): connection = await aio_pika.connect_robust( host=host, port=port, virtualhost=virtualhost, login=login, password=password ) channel = await connection.channel() # 关键:设置预取数,控制RabbitMQ一次推送到客户端的消息上限 await channel.set_qos(prefetch_count=BATCH_SIZE) queue = await channel.declare_queue(queue_name, auto_delete=True) await queue.bind(exchange=queue_ex_name, routing_key=queue_routing_key) await queue.consume(callback=process_message, consumer_tag='my_tag') async def process_message(msg): global pending_msg_count pending_msg_count += 1 try: # 原有消息处理逻辑 await msg.ack() finally: pending_msg_count -= 1 # 无待处理的已接收消息,触发通知 if pending_msg_count == 0: # 此处写入你给其他服务发通知的逻辑 send_notification()
方案2:直接读取内部缓冲状态
如果不想修改原有预取逻辑,也可以直接读取aio-pika的内部缓冲属性,注意这是私有属性,不同版本可能存在兼容性变化:
async def process_message(msg): # 原有消息处理逻辑 await msg.ack() # 获取当前消费者的内部缓冲队列长度 consumer = msg.channel._consumers[msg.consumer_tag] pending_buffer_size = consumer._queue.qsize() if pending_buffer_size == 0: # 无待处理的已接收消息,触发通知 send_notification()
注意:方案2依赖私有属性,跨版本兼容差,生产环境优先使用方案1。
内容的提问来源于stack exchange,提问作者Helgi Borg
相关产品推荐
相关产品推荐

