Python多线程实现:如何并行运行两个函数提升处理效率?
Python中GPU图像提取与IO保存并行的实现方案
最优方案是采用生产者-消费者模型,通过队列解耦GPU提取(生产者)和磁盘保存(消费者)两个任务,让它们并行执行。因为GPU提取属于计算密集型任务,而磁盘保存是IO密集型任务,这种模型能最大化利用CPU和GPU资源,避免串行执行时的等待浪费。
具体实现方式:用queue.Queue结合多线程
1. 核心思路
- 生产者线程:负责用GPU神经网络从视频中提取图像批次,将批次放入队列。
- 消费者线程:从队列中取出图像批次,写入磁盘。可以启动多个消费者线程,提升IO处理效率。
- 队列设置最大容量,防止GPU提取过快导致内存溢出。
2. 代码示例
import threading import queue import cv2 from pathlib import Path # 替换成你的实际GPU提取逻辑 def gpu_extract_batch(video_path, batch_size, output_queue): cap = cv2.VideoCapture(str(video_path)) frame_batch = [] frame_idx = 0 while cap.isOpened(): ret, frame = cap.read() if not ret: break frame_batch.append((frame_idx, frame)) frame_idx += 1 # 批次满了就放入队列 if len(frame_batch) == batch_size: output_queue.put(frame_batch) frame_batch = [] # 处理剩余的帧 if frame_batch: output_queue.put(frame_batch) # 放入结束标记,通知消费者任务完成 output_queue.put(None) cap.release() # 保存帧的消费者函数 def save_frame_batch(save_root, input_queue): save_root = Path(save_root) save_root.mkdir(exist_ok=True, parents=True) while True: batch = input_queue.get() # 收到结束标记则退出 if batch is None: input_queue.task_done() break # 保存当前批次的所有帧 for frame_idx, frame in batch: save_path = save_root / f"frame_{frame_idx}.jpg" cv2.imwrite(str(save_path), frame) input_queue.task_done() def main(): video_path = "your_video_dir/video.mp4" save_dir = "your_save_dir" batch_size = 32 # 队列最大容量设为2,避免内存占用过高 frame_queue = queue.Queue(maxsize=2) # 启动2个消费者线程(IO密集型任务适合多线程) consumer_threads = [] for _ in range(2): t = threading.Thread(target=save_frame_batch, args=(save_dir, frame_queue)) t.start() consumer_threads.append(t) # 启动生产者线程(GPU提取单线程即可,多线程无法提升GPU利用率) producer_thread = threading.Thread(target=gpu_extract_batch, args=(video_path, batch_size, frame_queue)) producer_thread.start() # 等待生产者完成所有提取任务 producer_thread.join() # 等待队列中所有批次处理完毕 frame_queue.join() # 给每个消费者发送结束标记 for _ in consumer_threads: frame_queue.put(None) # 等待所有消费者线程退出 for t in consumer_threads: t.join() if __name__ == "__main__": main()
3. 关键注意点
- GPU线程数量:GPU计算任务通常单线程就能占满GPU资源,无需启动多个生产者线程,反而可能引发框架上下文冲突。
- 队列容量:根据你的内存大小调整
maxsize,如果GPU提取远快于保存,队列太大容易导致内存被帧数据占满。 - 文件名冲突:保存时用全局帧索引或唯一标识命名,避免不同批次的帧文件名重复。
- 线程安全:如果你的GPU提取框架(如PyTorch/TensorFlow)有特定的线程要求,确保生产者线程是唯一调用GPU逻辑的线程。
替代方案:concurrent.futures.ThreadPoolExecutor
如果不想手动管理队列和线程,可以用线程池简化保存逻辑,让生产者提取完批次后直接提交保存任务到线程池:
from concurrent.futures import ThreadPoolExecutor import cv2 from pathlib import Path def gpu_extract_batch(video_path, batch_size): cap = cv2.VideoCapture(str(video_path)) frame_batch = [] frame_idx = 0 while cap.isOpened(): ret, frame = cap.read() if not ret: break frame_batch.append((frame_idx, frame)) frame_idx += 1 if len(frame_batch) == batch_size: yield frame_batch frame_batch = [] if frame_batch: yield frame_batch cap.release() def save_frame_batch(save_root, batch): save_root = Path(save_root) save_root.mkdir(exist_ok=True, parents=True) for frame_idx, frame in batch: save_path = save_root / f"frame_{frame_idx}.jpg" cv2.imwrite(str(save_path), frame) def main(): video_path = "your_video_dir/video.mp4" save_dir = "your_save_dir" batch_size = 32 # 启动线程池,设置2个线程处理保存任务 with ThreadPoolExecutor(max_workers=2) as executor: # 提取批次并提交保存任务 for batch in gpu_extract_batch(video_path, batch_size): executor.submit(save_frame_batch, save_dir, batch) if __name__ == "__main__": main()
这种方式代码更简洁,但队列的方式对任务流程的控制更精细,适合需要严格把控内存的场景。
内容的提问来源于stack exchange,提问作者MGeureka
相关产品推荐
相关产品推荐

