如何用Python多进程规避循环中的IO延迟问题?
用生产者-消费者架构规避IO开销的解决方案
问题根源:之前的多进程尝试为何更慢?
你之前的代码每次循环都创建一个新的Pool,进程的启动和销毁会带来巨大的开销,远远超过了IO操作本身的耗时,这是导致耗时从0.1秒暴涨到17秒的核心原因。多进程的优势只有在复用进程、避免频繁创建销毁时才能体现。
正确的生产者-消费者实现
针对你的场景,主进程需要专注于实时的图像预处理、模型预测和控制指令发送(CPU/GPU密集型),而把图像保存这种IO密集型任务交给独立的子进程处理,通过队列传递任务,完全不阻塞主循环。
代码实现
import multiprocessing from PIL import Image from io import BytesIO import base64 import numpy as np # 消费者函数:专门处理图像保存的IO操作 def image_saver_worker(queue): while True: # 从队列获取任务,阻塞直到有任务 img_c, steering_angle = queue.get() # 处理保存逻辑 save_image(img_c, steering_angle) # 标记任务完成,避免队列积压 queue.task_done() if __name__ == '__main__': # 创建任务队列,限制队列大小防止内存溢出(比如最多存100个待保存任务) save_queue = multiprocessing.Queue(maxsize=100) # 启动1-2个消费者进程(IO密集型任务不需要太多进程) for _ in range(2): worker = multiprocessing.Process(target=image_saver_worker, args=(save_queue,)) # 设置为守护进程,主进程退出时自动终止 worker.daemon = True worker.start() # 主循环:生产者逻辑 for data in data_arr: speed = float(data['speed']) image = Image.open(BytesIO(base64.b64decode(data['image']))) image = np.asarray(image) img_c = image.copy() image = img_preprocess(image) image = np.array([image]) # 模型预测(核心实时逻辑) steering_angle = float(model_steer.predict(image)) throttle = 1.0 - speed / speed_limit # 将保存任务放入队列,非阻塞,主循环继续执行 try: save_queue.put((img_c, steering_angle), block=False) except multiprocessing.Queue.Full: # 队列满时可选择丢弃旧任务或等待,这里简单打印提示 print("保存队列已满,跳过当前图像保存") print('{} {} {}'.format(steering_angle, throttle, speed)) send_control(steering_angle, throttle) # 等待所有队列任务处理完成(如果需要主进程等待保存完再退出) save_queue.join()
关键优化点
- 复用消费者进程:启动后持续运行,避免频繁创建销毁进程的开销
- 非阻塞任务投递:主进程调用
put(block=False)不会被IO操作阻塞,确保实时性 - 队列大小限制:防止过多未处理任务占用内存,队列满时可根据需求调整策略(丢弃/等待)
- 守护进程设置:主进程退出时消费者自动终止,无需手动管理
额外建议
- 如果你的
save_image函数涉及复杂的文件操作,可以考虑批量保存进一步优化,但单进程处理IO已经足够缓解主循环压力 - 确保
img_c在传递到队列时是可序列化的(numpy数组和PIL图像都支持multiprocessing队列的序列化)
内容的提问来源于stack exchange,提问作者Vikas Kumar Ojha
相关产品推荐
相关产品推荐

