Python多进程中如何串联IO密集与CPU密集操作?
串联多线程IO读取与多进程CPU计算的实现方案
你不需要必须先完成所有文件的读取再启动计算,有两种实用方案可以选择,分别对应不同的场景:
方案一:批量读取后批量计算(实现简单)
适合文件数量不多、单文件数据量较小的场景,先通过多线程一次性读完所有文件,再用多进程处理全部数据。代码示例如下:
from multiprocessing import Pool from multiprocessing.dummy import Pool as ThreadPool # IO密集型:读取文件并预处理数据 def read_file(file_path): with open(file_path, 'r') as f: # 根据你的实际数据格式调整预处理逻辑,这里假设每行是一个数值 data = [float(line.strip()) for line in f if line.strip()] return data # CPU密集型:数值计算逻辑 def process_data(data): # 替换为你的实际计算需求,示例为求和后取平方 return sum(data) ** 2 if __name__ == '__main__': # 替换为你的实际文件路径列表 file_paths = ["data_1.txt", "data_2.txt", "data_3.txt"] # 多线程读取:线程数可根据IO性能调整(比如设为CPU核心数的2-4倍) thread_pool = ThreadPool(4) all_raw_data = thread_pool.map(read_file, file_paths) thread_pool.close() thread_pool.join() # 多进程计算:进程数一般设为CPU核心数 process_pool = Pool(4) calculation_results = process_pool.map(process_data, all_raw_data) process_pool.close() process_pool.join() print("最终计算结果:", calculation_results)
方案二:流水线式异步处理(资源利用更高效)
如果文件数量多、单文件数据量大,批量读取会占用过多内存,这时可以用队列衔接多线程读取和多进程计算,让读取和计算并行进行——读好一个文件就立刻交给进程计算,不用等全部读取完成。代码示例如下:
from multiprocessing import Pool, Queue, Manager from multiprocessing.dummy import Pool as ThreadPool import threading # IO密集型:读取文件后放入队列 def read_to_queue(file_paths, data_queue): def _read_single(file_path): with open(file_path, 'r') as f: data = [float(line.strip()) for line in f if line.strip()] data_queue.put(data) thread_pool = ThreadPool(4) thread_pool.map(_read_single, file_paths) thread_pool.close() thread_pool.join() # 放入结束标记,通知进程没有更多数据 data_queue.put(None) # CPU密集型:从队列取数据并计算 def compute_from_queue(data_queue, result_list): while True: data = data_queue.get() if data is None: break # 替换为你的实际计算逻辑 result = sum(data) ** 2 result_list.append(result) if __name__ == '__main__': file_paths = ["data_1.txt", "data_2.txt", "data_3.txt", "data_4.txt"] # 用进程安全队列衔接读取和计算 data_queue = Queue() # 用Manager.List存储跨进程的计算结果 result_list = Manager().list() # 启动后台线程处理文件读取 read_thread = threading.Thread(target=read_to_queue, args=(file_paths, data_queue)) read_thread.start() # 启动多进程处理计算任务 process_pool = Pool(4) # 给每个进程分配计算任务 for _ in range(4): process_pool.apply_async(compute_from_queue, args=(data_queue, result_list)) process_pool.close() process_pool.join() read_thread.join() print("最终计算结果:", list(result_list))
两种方案的选择建议
- 批量处理:代码更简洁,调试方便,适合小体量数据场景。
- 流水线处理:内存占用更低,能同时利用IO和CPU资源,适合大数据量、多文件的场景,整体运行效率更高。
内容的提问来源于stack exchange,提问作者Aenaon
相关产品推荐
相关产品推荐

