You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)
AA3545624.7361fam1121
AA35456107.3580fam1121
AA3545672.0639fam1121
AA3545643.8766vic1121
AA354562382.8700vic1121
AA3544449.6488vic1121
AA3544472.0639fam3121
AA3544443.8766vic3121
AA3544472.0639fam3121
AA3544443.8766vic3121
AA3544472.0639fam3123
AA3544243.8766vic3123
AA3544472.0639fam6126
AA3544243.8766vic6126
AA3544472.0639fam6126
AA3544243.8766vic6126
AA3544472.0639fam6126
AA3544243.8766vic6126

CPU使用率情况

CPU未被充分利用(使用率维持在较低水平)


问题分析与解决方案

核心问题

  1. 多进程逻辑完全错误:p.map(calcIndex, DF_Feces_FAVF)会把DataFrame的每一行拆成单独任务,但你的calcIndex是处理分组的逻辑,子任务太小导致进程间通信开销远大于计算开销,CPU根本跑不满。
  2. 计算逻辑冗余:多次重复groupby.apply,完全没用到pandas的向量化优势,效率极低。
  3. 进程调用顺序混乱: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 01:35:25