Python多进程任务仅获2倍提速,求8核CPU下优化方案
优化多进程资源利用的实用方法
你的8核CPU仅获得2倍提速,核心原因大概率是IO密集环节(读、写)与CPU密集环节(scipy处理)未解耦,导致CPU资源被IO等待浪费。以下是针对性优化方案:
1. 解耦IO与CPU任务,采用生产者-消费者模型
将读数据、处理数据、写数据拆分为独立阶段,用队列衔接不同任务,让IO操作和CPU计算并行执行,避免互相等待。
示例代码:
import concurrent.futures from queue import Queue import pandas as pd # 假设read返回原始数据,treat接收原始数据并返回处理后的DataFrame def read(time): # 实际读取逻辑 return raw_data def treat(time, raw_data): # 实际scipy处理逻辑 return pd.DataFrame(processed_data) def read_worker(schedule_queue, data_queue): while not schedule_queue.empty(): time = schedule_queue.get() raw_data = read(time) data_queue.put((time, raw_data)) schedule_queue.task_done() def treat_worker(data_queue, result_queue): while True: time, raw_data = data_queue.get() if time is None: # 任务结束信号 data_queue.task_done() break df = treat(time, raw_data) result_queue.put((time, df)) data_queue.task_done() def write_worker(result_queue): while True: time, df = result_queue.get() if time is None: # 任务结束信号 result_queue.task_done() break df.to_csv(f'partial_{time}_data.csv') result_queue.task_done() if __name__ == "__main__": schedule = [1,2,3,4,5,6,7,8] # 示例时间列表 schedule_queue = Queue() data_queue = Queue(maxsize=4) # 限制队列大小,避免内存溢出 result_queue = Queue(maxsize=4) # 填充任务队列 for time in schedule: schedule_queue.put(time) # 启动IO读进程(2个,适配IO密集型任务) with concurrent.futures.ProcessPoolExecutor(max_workers=2) as read_exec: read_exec.submit(read_worker, schedule_queue, data_queue) # 启动CPU处理进程(8个,匹配核心数) with concurrent.futures.ProcessPoolExecutor(max_workers=8) as treat_exec: treat_workers = [treat_exec.submit(treat_worker, data_queue, result_queue) for _ in range(8)] # 启动IO写进程(2个) with concurrent.futures.ProcessPoolExecutor(max_workers=2) as write_exec: write_worker_fut = write_exec.submit(write_worker, result_queue) # 等待读任务完成 schedule_queue.join() # 发送结束信号给处理进程 for _ in range(8): data_queue.put((None, None)) data_queue.join() # 发送结束信号给写进程 result_queue.put((None, None)) result_queue.join()
2. 调整进程池大小,匹配任务类型
- CPU密集型任务(treat):
max_workers设置为CPU核心数(8)即可,过多进程会导致上下文切换开销激增。 - IO密集型任务(read/write):可设置为核心数的2-4倍,但不要超过10,避免IO资源竞争。
你之前设置max_workers=16,对于8核CPU来说会加重上下文切换负担,建议先将process_map的max_workers改为8测试效果。
3. 优化CPU密集的treat函数
先提升单任务处理速度,再放大多进程收益:
- 用
numba对核心计算逻辑做JIT编译,加速数值运算。 - 检查scipy函数是否有内置并行参数(如
scipy.optimize的workers),开启内置并行。 - 用numpy向量化操作替代Python循环,减少解释层开销。
示例(numba加速):
from numba import jit @jit(nopython=True) def core_calculation(data): # 这里是scipy处理中的核心数值计算逻辑 result = ... return result def treat(time, raw_data): processed = core_calculation(raw_data) # 其他scipy操作或转换为DataFrame df = pd.DataFrame(processed) return df
4. 优化IO操作效率
- 读取数据时,优先用Parquet、Feather等二进制格式替代文本格式;若必须用CSV,使用
pandas.read_csv(engine='pyarrow')加速读取。 - 写入CSV时,用
df.to_csv(engine='pyarrow'),比默认引擎快2-3倍;允许的话先写入Parquet,后续再批量转CSV,IO效率会大幅提升。 - 批量处理IO:合并多个time的数据一次性读取,处理后再拆分写入,减少IO调用次数。
5. 避免不必要的进程间数据拷贝
若read读取的数据需要在进程间传递,使用multiprocessing.Array或multiprocessing.Manager的共享容器,避免完整数据拷贝。如果每个子进程独立处理自身time的数据(如当前thread函数逻辑),则无需额外处理,但要确保无隐式跨进程数据传递。
内容的提问来源于stack exchange,提问作者jAke
相关产品推荐
相关产品推荐

