Jupyter Notebook中multiprocessing未充分利用CPU核心的求助
问题描述
我用的是6物理核心、12逻辑线程的i7处理器,现在要给一个120万行的pandas DataFrame(只取了1/3子集)计算香农多样性指数。但子集任务跑了快10小时还没结束,CPU也没跑满,环境是Jupyter Notebook,相关信息如下:
运行代码
import pandas as pd import matplotlib as plt import numpy as np import seaborn as sns from support_functions import functions from multiprocessing import Pool p= Pool(12) DF = pd.read_csv("H:/data_micro/Feces2tpsF.csv") #calculate the shannon diversity def calcIndex(DATAF, levelDF): DATAF["Pi"] = DATAF["log2Smooth"].groupby(level=["PatientID", "Timepoint"]).apply(lambda x: x / x.sum()) DATAF["pi*lnpi"]= DATAF["Pi"]*np.log(DATAF["Pi"]) DATAF["H_log2"] = DATAF["pi*lnpi"].groupby(level=["PatientID", "Timepoint"]).apply( lambda x: x.sum()*-1) return DATAF DF_Feces_FAVF = DF[(DF["Phylum"]== "vic")] DF_Feces_FAVF = DF_Feces_FAVF.set_index(["PatientID", "Timepoint", "Phylum"]) DF_Feces_FAVF = p.map(calcIndex,DF_Feces_FAVF ) p.start() p.close() p.join()
DataFrame结构
| 样本编号(Sample nr) | 平滑值(smooth) | Phylum | 时间点(timepoint) | 患者ID(PatientID) |
|---|---|---|---|---|
| AA35456 | 24.7361 | fam | 1 | 121 |
| AA35456 | 107.3580 | fam | 1 | 121 |
| AA35456 | 72.0639 | fam | 1 | 121 |
| AA35456 | 43.8766 | vic | 1 | 121 |
| AA35456 | 2382.8700 | vic | 1 | 121 |
| AA35444 | 49.6488 | vic | 1 | 121 |
| AA35444 | 72.0639 | fam | 3 | 121 |
| AA35444 | 43.8766 | vic | 3 | 121 |
| AA35444 | 72.0639 | fam | 3 | 121 |
| AA35444 | 43.8766 | vic | 3 | 121 |
| AA35444 | 72.0639 | fam | 3 | 123 |
| AA35442 | 43.8766 | vic | 3 | 123 |
| AA35444 | 72.0639 | fam | 6 | 126 |
| AA35442 | 43.8766 | vic | 6 | 126 |
| AA35444 | 72.0639 | fam | 6 | 126 |
| AA35442 | 43.8766 | vic | 6 | 126 |
| AA35444 | 72.0639 | fam | 6 | 126 |
| AA35442 | 43.8766 | vic | 6 | 126 |
CPU使用率情况
CPU未被充分利用(使用率维持在较低水平)
问题分析与解决方案
核心问题
- 多进程逻辑完全错误:
p.map(calcIndex, DF_Feces_FAVF)会把DataFrame的每一行拆成单独任务,但你的calcIndex是处理分组的逻辑,子任务太小导致进程间通信开销远大于计算开销,CPU根本跑不满。 - 计算逻辑冗余:多次重复
groupby.apply,完全没用到pandas的向量化优势,效率极低。 - 进程调用顺序混乱:
p.start()应该在map前调用,且map本身会阻塞到任务完成,后续的start/close/join属于无效操作。
修复步骤
1. 重构计算逻辑(优先推荐,无需多进程)
用pandas向量化操作替代重复分组,一次完成计算,效率至少提升10倍:
import pandas as pd import numpy as np DF = pd.read_csv("H:/data_micro/Feces2tpsF.csv") DF_Feces_FAVF = DF[DF["Phylum"] == "vic"] # 一次分组获取各组总和,用transform广播到每行 group_key = ["PatientID", "Timepoint"] grouped = DF_Feces_FAVF.groupby(group_key)["log2Smooth"] total = grouped.transform("sum") # 计算Pi和香农指数 DF_Feces_FAVF["Pi"] = DF_Feces_FAVF["log2Smooth"] / total DF_Feces_FAVF["pi*lnpi"] = DF_Feces_FAVF["Pi"] * np.log(DF_Feces_FAVF["Pi"]) # 用transform把组内计算的香农指数广播到每行 DF_Feces_FAVF["H_log2"] = grouped.transform(lambda x: -(x / x.sum() * np.log(x / x.sum())).sum())
2. 正确使用多进程(仅当数据量极大时需要)
如果必须并行,要按分组拆分任务,而非逐行拆分,确保每个进程的计算量足够大:
import pandas as pd import numpy as np from multiprocessing import Pool def calc_single_group(group): total = group["log2Smooth"].sum() group["Pi"] = group["log2Smooth"] / total group["pi*lnpi"] = group["Pi"] * np.log(group["Pi"]) group["H_log2"] = -(group["pi*lnpi"].sum()) return group DF = pd.read_csv("H:/data_micro/Feces2tpsF.csv") DF_Feces_FAVF = DF[DF["Phylum"] == "vic"] # 按(PatientID, Timepoint)拆分分组 groups = [g for _, g in DF_Feces_FAVF.groupby(["PatientID", "Timepoint"])] # 用物理核心数初始化进程池(6个比12个更高效,避免超线程调度开销) with Pool(6) as p: result_groups = p.map(calc_single_group, groups) # 合并结果 DF_Feces_FAVF = pd.concat(result_groups)
3. 其他优化建议
- 读取CSV时指定列类型减少内存:比如
PatientID设为int32,Phylum设为category,示例:pd.read_csv("xxx.csv", dtype={"PatientID": np.int32, "Phylum": "category"}) - 把CSV转成Parquet格式,后续读取速度提升数倍:
DF.to_parquet("data.parquet"),读取用pd.read_parquet - Jupyter中用
%timeit测试单组计算耗时,快速定位瓶颈
内容的提问来源于stack exchange,提问作者Hedayat
相关产品推荐
相关产品推荐

