如何基于OpenVINO异步推理处理RabbitMQ队列中的图片?
OpenVINO异步推理适配图片列表/RabbitMQ输入方案
核心思路
OpenVINO异步API的核心是通过双请求交替执行(curr/next),将图片加载预处理和推理计算并行,提升吞吐量。不管输入是视频流还是单张图片,只要按顺序提供预处理后的图片张量,就能适配异步模式。
一、图片列表处理方案
直接遍历图片列表,维护两个异步推理请求,交替启动和等待,实现加载与推理并行。
示例代码
import openvino as ov import cv2 # 初始化模型与推理请求 core = ov.Core() model = core.read_model(model="your_model.xml") compiled_model = core.compile_model(model=model, device_name="CPU") input_layer = compiled_model.input(0) output_layer = compiled_model.output(0) # 创建两个异步推理请求 curr_request = compiled_model.create_infer_request() next_request = compiled_model.create_infer_request() # 待处理的图片路径列表 image_paths = ["img1.jpg", "img2.jpg", "img3.jpg"] # 预处理函数(根据模型要求调整,比如尺寸、归一化、通道顺序等) def preprocess(img_path): img = cv2.imread(img_path) resized_img = cv2.resize(img, (input_layer.shape[3], input_layer.shape[2])) # 若模型需要BGR转RGB、归一化等,在此添加处理 return ov.Tensor(resized_img) # 启动第一个请求 if image_paths: first_tensor = preprocess(image_paths[0]) curr_request.set_tensor(input_layer, first_tensor) curr_request.start_async() # 遍历剩余图片,交替执行请求 for idx, img_path in enumerate(image_paths[1:], start=1): # 准备下一张图的张量,启动next请求 next_tensor = preprocess(img_path) next_request.set_tensor(input_layer, next_tensor) next_request.start_async() # 等待当前请求完成,处理结果 curr_request.wait() result = curr_request.get_tensor(output_layer).data # 替换为你的结果逻辑:比如保存、打印、推送等 print(f"处理完成 {image_paths[idx-1]},结果形状: {result.shape}") # 交换curr和next请求,进入下一轮 curr_request, next_request = next_request, curr_request # 处理最后一个请求 if image_paths: curr_request.wait() last_result = curr_request.get_tensor(output_layer).data print(f"处理完成最后一张图 {image_paths[-1]},结果形状: {last_result.shape}")
关键点
- 预处理逻辑必须严格匹配模型要求(尺寸、通道顺序、归一化参数等)。
- 交替请求的方式让图片加载/预处理与推理并行,避免串行等待的时间损耗。
二、RabbitMQ消费者方案
RabbitMQ是消息驱动的输入源,需结合异步消费框架(如pika的异步模式),维护两个推理请求,收到消息后分配给空闲请求执行,实现高吞吐量。
示例代码
import openvino as ov import cv2 import numpy as np import pika from pika.adapters.asyncio_connection import AsyncioConnection import asyncio # 初始化OpenVINO模型 core = ov.Core() model = core.read_model(model="your_model.xml") compiled_model = core.compile_model(model=model, device_name="CPU") input_layer = compiled_model.input(0) output_layer = compiled_model.output(0) # 维护两个异步推理请求 requests = [compiled_model.create_infer_request(), compiled_model.create_infer_request()] active_idx = 0 # 当前正在运行的请求索引 # 图片预处理:从RabbitMQ字节数据解码并转换为模型所需张量 def preprocess_image(img_bytes): nparr = np.frombuffer(img_bytes, np.uint8) img = cv2.imdecode(nparr, cv2.IMREAD_COLOR) resized_img = cv2.resize(img, (input_layer.shape[3], input_layer.shape[2])) # 根据模型需求添加通道转换、归一化等操作 return ov.Tensor(resized_img) # 推理完成回调:处理结果并确认RabbitMQ消息 def infer_callback(request, channel, delivery_tag): result = request.get_tensor(output_layer).data # 替换为你的结果处理逻辑 print(f"处理完成图片,结果形状: {result.shape}") # 确认消息已处理,避免重复消费 channel.basic_ack(delivery_tag=delivery_tag) # 接收RabbitMQ消息的回调 async def on_message(ch, method, properties, body): global active_idx # 获取空闲请求索引 idle_idx = 1 - active_idx idle_request = requests[idle_idx] # 预处理图片并设置张量 input_tensor = preprocess_image(body) idle_request.set_tensor(input_layer, input_tensor) # 设置推理完成回调,绑定通道和消息标签 idle_request.set_callback(lambda req: infer_callback(req, ch, method.delivery_tag)) # 启动异步推理 idle_request.start_async() # 切换活跃请求索引 active_idx = idle_idx async def main(): # 初始化RabbitMQ异步连接 connection = AsyncioConnection(pika.ConnectionParameters(host="localhost")) channel = await connection.channel() await channel.queue_declare(queue="image_queue") # 设置预取数为2,匹配双请求的并行能力 await channel.basic_qos(prefetch_count=2) # 开始消费消息 await channel.basic_consume(queue="image_queue", on_message_callback=on_message) print("等待RabbitMQ图片消息,按CTRL+C退出") await asyncio.Future() # 保持服务运行 if __name__ == "__main__": asyncio.run(main())
关键点
- 使用RabbitMQ异步消费框架,避免阻塞等待消息。
- 通过
prefetch_count=2控制预取消息数,确保两个推理请求始终有任务可执行。 - 推理完成后必须确认消息,防止消息重复投递。
内容的提问来源于stack exchange,提问作者Ksenia Belozerova
相关产品推荐
相关产品推荐

