单函数内CPU利用vs多进程并行:公式计算方案选型与优化咨询
问题背景
我有数千个公式需要计算评估,例如rank(sqrt(v1))或corr(v1, square(v2))。所有变量(v1、v2……)已通过共享内存供多进程使用,所有计算均依赖numpy或scipy函数。
方案1:多进程并行计算每个公式
from scipy.stats import rankdata import multiprocessing as mp def myeval(formula): ... # 从共享内存读取数据 ... # 计算 return value def eval_access(formula): factor = myeval(formula) # 占用1GB内存 factor_performance1 = rankdata(factor, axis=1) # 再占用1GB内存 factor_performance2 = np.full_like(factor, np.nan) # 再占用1GB内存 ... # 其他计算 pickle_part # 记录性能的I/O操作(仅存储部分浮点数) return with mp.Pool(30) as p: p.map(eval_access, all_formulas)
该方案可充分利用CPU(利用率接近100%),但每个进程最多占用3GB内存,30个进程同时运行时最高可能占用60GB内存;且numpy操作在多进程环境下似乎会变慢。
方案2:单公式内多进程拆分计算
from scipy.stats import rankdata import multiprocessing as mp def myeval(ind1, ind2): # 修改value[ind1:ind2, :]或value[:, ind1:ind2] def myeval_multiprocess(formula): # 从共享内存读取数据 value = np.zeros(shape) # 准备返回值存储空间 ... # 共享value内存 ...# 将value的索引拆分为30份 with mp.Pool(30) as p: p.map(myeval, ind_list) return value def eval_access(formula): factor = myeval_multiprocess(formula) # 占用1GB内存 factor_performance1 = rankdata(factor, axis=1) # 再占用1GB内存 factor_performance2 = np.full_like(factor, np.nan) # 再占用1GB内存 ... # 其他计算 pickle_part # 记录性能的I/O操作(仅存储部分浮点数) return for formula in all_formulas: eval_access(formula)
理论上两种方案耗时相同,且后者内存占用更低,但实际存在两个问题:一是在myeval_multiprocess中共享内存耗时较长;二是eval_access内并非所有操作都能充分利用CPU,除非改造所有用到的函数以支持拆分计算。
咨询问题
- 哪种方案更优?
- 如何优化提速?
- 第二种方案能否达到理想状态?
- 是否存在简便的全计算多进程实现方法?
方案对比与优化建议
1. 哪种方案更优?
没有绝对最优,取决于你的资源瓶颈:
- 如果内存充足(能承载60GB峰值),方案1更省心,无需改造现有计算逻辑,CPU利用率拉满,适合快速落地。
- 如果内存紧张,或者单公式计算量极大(远超单进程算力),方案2是必要选择,但需要承担额外的代码改造成本。
2. 优化提速方法
针对方案1的优化:
- 合理设置进程数:numpy/scipy多数操作依赖OpenBLAS/MKL等多线程后端,30个进程+多线程会导致CPU上下文切换过载,反而变慢。建议进程数设为物理核心数的1/2到1倍(比如8核CPU开8个进程)。
- 主动释放内存:在
eval_access中,计算完不再使用的变量后,手动执行del factor+gc.collect(),减少单进程内存占用,降低整体内存峰值。 - 替换低效I/O:pickle对数值数据的存储效率不高,换成
numpy.savez或feather这类专用格式,减少I/O耗时。 - 批量相似公式:对逻辑相近的公式做批量处理,减少进程创建销毁的重复开销。
针对方案2的优化:
- 复用共享内存:提前初始化全局共享内存池,用
multiprocessing.Array配合numpy.frombuffer直接映射,避免每次计算重复创建共享内存的开销。 - 用
imap_unordered替代map:若计算结果顺序不影响后续逻辑,imap_unordered能在子进程完成任务后立即返回结果,减少等待时间。 - 对齐内存块拆分:拆分索引时尽量按数组连续内存块(比如行维度)拆分,避免非连续切片导致的numpy性能损耗。
- 用numba加速计算:如果
myeval内的逻辑可编译,用numba能大幅提升单线程效率,抵消多进程的额外开销。
3. 第二种方案能否达到理想状态?
可以,但需要解决两个核心问题:
- 共享内存开销:通过预分配全局共享内存、直接映射内存的方式,能把共享内存的耗时降到可忽略的程度。
- 全流程并行化:对于
rankdata这类无法直接拆分的函数,可手动实现并行版本——比如按行拆分factor,子进程分别处理部分行的排名后合并;或者用dask.array这类并行数组库,它会自动拆分计算任务,无需手动写多进程逻辑。
4. 简便的全计算多进程实现方法
推荐用Dask或joblib替代手动多进程,代码改动小且兼顾效率:
方案A:Dask自动并行
Dask能自动将numpy/scipy逻辑拆分为并行任务,支持共享内存,内存占用可控:
import dask.array as da from scipy.stats import rankdata from dask import delayed, compute def myeval_dask(formula): # 从共享内存读取数据转为dask array(按行拆分块) v1 = da.from_array(shared_v1, chunks=(1000, -1)) # 解析公式计算(示例:rank(sqrt(v1))) result = da.map_blocks(lambda x: rankdata(x, axis=1), da.sqrt(v1), dtype=float) return result.compute() def eval_access(formula): factor = myeval_dask(formula) # 后续计算逻辑不变 ... # 并行处理所有公式 tasks = [delayed(eval_access)(formula) for formula in all_formulas] compute(*tasks, num_workers=8)
方案B:joblib简化多进程
joblib的Parallel模块支持多进程/多线程切换,自动处理numpy数组的内存共享:
from joblib import Parallel, delayed Parallel(n_jobs=8, backend='loky')(delayed(eval_access)(formula) for formula in all_formulas)
loky后端能避免多进程下的numpy数据拷贝问题,比原生multiprocessing更高效。
内容的提问来源于stack exchange,提问作者LeonM
相关产品推荐
相关产品推荐

