如何在Python嵌套循环中并行计算以生成距离矩阵?
如何通过并行计算加速嵌套循环的距离矩阵计算?
当然可以!你的场景完全适合并行化——计算任务彼此独立、单任务内存占用极低,并行化能直接利用多核CPU大幅提升速度。下面分步骤给出优化方案:
先做串行优化(必做,不并行也能提速)
原代码存在重复计算:同一个smi会被lang.convert转换len(smis)次。先提前批量转换所有特征,避免冗余操作:
import numpy as np for n in id: df_ = df[df['id_label'] == n] smis = df_['origin'].tolist() # 提前批量转换所有smis到特征,避免重复计算 features = [lang.convert(smi) for smi in smis] distance_matrix = [] for f1 in features: row = [get_distance(f1, f2) for f2 in features] distance_matrix.append(row) distance_matrix = np.array(distance_matrix) df_["distances"] = distance_matrix.sum(axis=1) df_.to_csv(f"output_cluster_{n}.csv", index=False)
并行化方案1:按cluster并行(实现简单,见效快)
不同id_label对应的cluster处理完全独立,直接把每个cluster的计算任务放到多进程中执行:
import numpy as np from concurrent.futures import ProcessPoolExecutor def process_cluster(n, df): df_ = df[df['id_label'] == n].copy() # 拷贝避免修改原DataFrame smis = df_['origin'].tolist() features = [lang.convert(smi) for smi in smis] distance_matrix = [] for f1 in features: row = [get_distance(f1, f2) for f2 in features] distance_matrix.append(row) distance_matrix = np.array(distance_matrix) df_["distances"] = distance_matrix.sum(axis=1) df_.to_csv(f"output_cluster_{n}.csv", index=False) return f"Cluster {n} processed" # 主程序入口(Windows系统必须加这个判断) if __name__ == "__main__": # 自动适配CPU核心数,也可以手动指定max_workers with ProcessPoolExecutor() as executor: # 提交所有cluster的计算任务 results = executor.map(process_cluster, id, [df]*len(id)) # 可选:打印任务完成状态 for res in results: print(res)
并行化方案2:cluster内按行并行(适合单个cluster数据量大的场景)
如果某个cluster的smi数量极多,还可以把每行距离的计算也并行化:
import numpy as np from concurrent.futures import ProcessPoolExecutor def compute_row(f1, features): # 单一行的距离计算逻辑 return [get_distance(f1, f2) for f2 in features] def process_cluster(n, df): df_ = df[df['id_label'] == n].copy() smis = df_['origin'].tolist() features = [lang.convert(smi) for smi in smis] # 并行计算每一行的距离 with ProcessPoolExecutor() as executor: distance_matrix = list(executor.map(compute_row, features, [features]*len(features))) distance_matrix = np.array(distance_matrix) df_["distances"] = distance_matrix.sum(axis=1) df_.to_csv(f"output_cluster_{n}.csv", index=False) return f"Cluster {n} processed" if __name__ == "__main__": with ProcessPoolExecutor() as executor: results = executor.map(process_cluster, id, [df]*len(id)) for res in results: print(res)
关键注意事项
- 用
ProcessPoolExecutor而非ThreadPoolExecutor:因为lang.convert和get_distance属于CPU密集型任务,Python的GIL会限制线程并行效率,多进程能绕开GIL充分利用多核。 - Windows系统必须把主逻辑放在
if __name__ == "__main__":下,否则会触发进程启动错误。 - 若
lang.convert支持批量输入(比如直接传列表),优先用批量接口,进一步减少转换耗时。 - 可通过
max_workers参数调整进程数,比如ProcessPoolExecutor(max_workers=6),一般不超过CPU核心数的1.5倍即可。
内容的提问来源于stack exchange,提问作者Giano
相关产品推荐
相关产品推荐

