如何对50行22000列的pandas DataFrame并行计算列间距离相关性(dcor)
并行计算DataFrame列对间距离相关性的最优方案
看你的问题,是要处理一个50行、22000列的DataFrame,用dcor计算所有列对的距离相关性——串行肯定慢到离谱,毕竟22000列要算差不多2.4亿个列对(准确说是(22000*22001)/2≈2.42亿),串行跑不知道要等到猴年马月。下面给你几个实用的并行实现方案,从最省心到最灵活的都有:
方案一:用dcor官方的并行函数(最推荐)
其实dcor本身就提供了批量计算距离相关矩阵的并行接口,完全不用自己写循环或者管理进程,代码简洁还高效。直接用dcor.pairwise.pairwise_distance_correlation函数,它支持n_jobs参数,可以指定用多少个CPU核心:
import pandas as pd import dcor import numpy as np # 替换成你的真实DataFrame np.random.seed(42) DF = pd.DataFrame(np.random.randn(50, 22000)) # 调用官方并行函数,n_jobs=-1表示用所有可用CPU核心 dc_correlation_matrix = dcor.pairwise.pairwise_distance_correlation(DF.values, n_jobs=-1) # 转成DataFrame格式,保持行列索引和原数据一致 DCOR_REZ = pd.DataFrame(dc_correlation_matrix, index=DF.columns, columns=DF.columns)
为什么推荐这个?
- 官方实现,经过优化,比自己写的并行逻辑效率更高
- 代码量极少,不用手动处理列对生成、结果填充这些繁琐步骤
- 自动处理对称矩阵的对称性,不用你手动赋值
DCOR_REZ.loc[j,i] = dc
方案二:用joblib实现自定义并行(灵活可控)
如果需要在计算距离相关性前后加一些自定义逻辑(比如对数据做预处理、过滤某些列对),可以用joblib的Parallel和delayed来实现并行:
import pandas as pd import dcor from joblib import Parallel, delayed import numpy as np from tqdm import tqdm # 可选,用来显示进度条 # 替换成你的真实DataFrame np.random.seed(42) DF = pd.DataFrame(np.random.randn(50, 22000)) col_names = DF.columns.tolist() n_cols = len(col_names) # 初始化结果矩阵,对角线初始化为1(自己和自己的距离相关性是1) DCOR_REZ = pd.DataFrame(np.eye(n_cols), index=col_names, columns=col_names) # 生成所有需要计算的列对(只算上三角,i < j,避免重复计算) def generate_column_pairs(): pairs = [] for i in range(n_cols): for j in range(i + 1, n_cols): pairs.append((i, j)) return pairs # 定义单个列对的计算函数 def compute_single_dcor(i, j): # 转成numpy数组,比pandas Series处理更快 v1 = DF.iloc[:, i].values v2 = DF.iloc[:, j].values return i, j, dcor.distance_correlation(v1, v2) # 并行计算,n_jobs=-1用所有CPU核心,verbose=10可以看到进度 pairs = generate_column_pairs() results = Parallel(n_jobs=-1, verbose=10)( delayed(compute_single_dcor)(i, j) for i, j in tqdm(pairs) ) # 把结果填充到矩阵里 for i, j, dc_value in results: DCOR_REZ.iloc[i, j] = dc_value DCOR_REZ.iloc[j, i] = dc_value
注意事项:
- 用
iloc代替loc处理整数索引,速度更快 - 把数据转成numpy数组传递给dcor,减少pandas Series的额外开销
- 用
tqdm可以直观看到计算进度,避免不知道程序跑到哪了
方案三:用multiprocessing手动管理进程池(底层可控)
如果你需要更精细地控制进程的创建和调度(比如限制进程数、处理异常),可以用Python标准库的multiprocessing模块:
import pandas as pd import dcor import multiprocessing as mp import numpy as np from tqdm import tqdm # 替换成你的真实DataFrame np.random.seed(42) DF = pd.DataFrame(np.random.randn(50, 22000)) col_names = DF.columns.tolist() n_cols = len(col_names) DCOR_REZ = pd.DataFrame(np.eye(n_cols), index=col_names, columns=col_names) # 生成列对 column_pairs = [(i, j) for i in range(n_cols) for j in range(i + 1, n_cols)] # 单个列对的计算函数 def compute_dcor(pair): i, j = pair v1 = DF.iloc[:, i].values v2 = DF.iloc[:, j].values return i, j, dcor.distance_correlation(v1, v2) # Windows系统必须把主逻辑放在if __name__ == "__main__"里,避免进程启动问题 if __name__ == "__main__": # 创建进程池,用所有CPU核心 with mp.Pool(processes=mp.cpu_count()) as pool: # 用tqdm显示进度 results = list(tqdm(pool.imap(compute_dcor, column_pairs), total=len(column_pairs))) # 填充结果矩阵 for i, j, dc_value in results: DCOR_REZ.iloc[i, j] = dc_value DCOR_REZ.iloc[j, i] = dc_value
注意事项:
- Windows系统下必须加
if __name__ == "__main__",否则会出现进程重复启动的问题 pool.imap比pool.map更适合显示进度,因为它是迭代返回结果的
关键注意点(必看)
- 内存问题:22000x22000的float64矩阵大概占3.7GB内存,确保你的机器有足够的可用内存。如果内存不足,可以考虑:
- 用
float32类型存储结果(会损失一点精度,需要权衡) - 分块计算,每次只处理一部分列对,分批填充结果
- 用
- 数据类型:尽量用numpy数组而不是pandas Series传递给dcor,能显著提升计算速度
- 进度监控:加上
tqdm或者joblib的verbose参数,避免长时间等待不知道程序状态
内容的提问来源于stack exchange,提问作者eod
相关产品推荐
相关产品推荐

