如何为嵌套循环内循环实现多进程,避免进程池重复创建开销
问题描述
我在实现一篇论文中的基于归一化压缩距离(NCD)的分类算法,初始代码在小规模数据下可正常运行,但数据集规模增大(实际为数十万长度的核苷酸序列)后,计算时间急剧增加。尝试用multiprocessing优化内循环的NCD计算时,每次处理测试样本都创建新进程池,导致进程创建销毁的开销过大,反而比单进程更慢。需要优化内循环的并行计算,同时避免重复管理进程池的开销。
初始单进程代码:
import gzip import numpy as np k = 16 training_set = np.array([('TCCCTACACT', 9), ('AGTTGGTATT', 12), ('AGTGGATCAC', 8), ('CGCAAGTGTG', 3), ('CTGTTCCCCC', 7), ('CGACTGGTAA', 10), ('CGCCGAGAAG', 4), ('TAGCTACGAC', 5), ('CCTTGCGCGT', 11), ('CCAAAAGAAA', 12), ('CCTCAGGAGG', 6), ('AGGCCACTTA', 5), ('GCGGGAACGG', 5), ('CTATTACCAA', 2), ('ACACTTTTTT', 8), ('GATGCAGCGT', 1)]) test_set= np.array([('GATGCAGCGT', 3), ('GATGCAGCGT', 0), ('GATGCAGCGT', 3), ('GATGCAGCGT', 4)]) for ( x1, _ ) in test_set: Cx1 = len( gzip . compress ( x1.encode() ) ) distance_from_x1 = [] for ( x2, _ ) in training_set: Cx2 = len( gzip.compress(x2.encode())) x1x2 = " ".join([ x1 , x2 ]) Cx1x2 = len( gzip.compress ( x1x2.encode() )) ncd = ( Cx1x2 - min( Cx1 , Cx2 ) ) / max(Cx1 , Cx2 ) distance_from_x1.append( ncd ) sorted_idx = np.argsort ( np.array(distance_from_x1 ) ) top_k_class = list(training_set [ sorted_idx[: k] , 1]) predict_class = max(set( top_k_class ) ,key = top_k_class.count ) print(predict_class)
尝试的多进程代码(存在重复创建进程池的问题):
import gzip import numpy as np from multiprocessing import Pool, cpu_count def calculate_ncd(x1, x2): Cx1 = len(gzip.compress(x1.encode())) Cx2 = len(gzip.compress(x2.encode())) x1x2 = " ".join([x1, x2]) Cx1x2 = len(gzip.compress(x1x2.encode())) ncd = (Cx1x2 - min(Cx1, Cx2)) / max(Cx1, Cx2) return ncd def process_train_sample(x1, x2): return calculate_ncd(x1, x2) def process_test_sample(test_sample): x1, _ = test_sample pool = Pool(cpu_count()) distance_from_x1 = pool.starmap(process_train_sample, [([x1] * len(training_set), training_set)]) pool.close() pool.join() sorted_idx = np.argsort(np.array(distance_from_x1)) top_k_class = list(training_set[sorted_idx[:k], 1]) predict_class = max(set(top_k_class), key=top_k_class.count) return predict_class k = 16 training_set = np.array([('TCCCTACACT', 9), ('AGTTGGTATT', 12), ('AGTGGATCAC', 8), ('CGCAAGTGTG', 3), ('CTGTTCCCCC', 7), ('CGACTGGTAA', 10), ('CGCCGAGAAG', 4), ('TAGCTACGAC', 5), ('CCTTGCGCGT', 11), ('CCAAAAGAAA', 12), ('CCTCAGGAGG', 6), ('AGGCCACTTA', 5), ('GCGGGAACGG', 5), ('CTATTACCAA', 2), ('ACACTTTTTT', 8), ('GATGCAGCGT', 1)]) test_set= np.array([('GATGCAGCGT', 3), ('GATGCAGCGT', 0), ('GATGCAGCGT', 3), ('GATGCAGCGT', 4)]) results = [process_test_sample(sample) for sample in test_set]# Process test samples in parallel print(results) # Print the predicted classes
解决方案
核心优化点:
- 全局复用进程池:只创建一次进程池,所有测试样本的内循环计算都复用该池,避免反复创建销毁进程的开销
- 预计算训练集压缩长度:训练集固定,提前计算所有样本的压缩长度,避免每次计算NCD时重复压缩训练样本
- 减少重复计算:每个测试样本的压缩长度只计算一次,无需在每个进程任务中重复计算
修改后的代码:
import gzip import numpy as np from multiprocessing import Pool, cpu_count # 预计算训练集样本的压缩长度,避免重复计算 def precompute_training_cx(training_set): training_data = [] for seq, label in training_set: cx = len(gzip.compress(seq.encode())) training_data.append((seq, label, cx)) return training_data # 优化后的NCD计算,直接传入预计算好的压缩长度 def calculate_ncd(x1, x2, cx1, cx2): x1x2 = " ".join([x1, x2]) cx1x2 = len(gzip.compress(x1x2.encode())) return (cx1x2 - min(cx1, cx2)) / max(cx1, cx2) def process_test_sample(test_sample, training_data, pool): x1, _ = test_sample # 仅计算一次当前测试样本的压缩长度 cx1 = len(gzip.compress(x1.encode())) # 构造所有并行任务:配对测试样本与每个训练样本,带上预计算的压缩长度 tasks = [(x1, seq, cx1, cx2) for seq, _, cx2 in training_data] # 用全局进程池并行计算所有NCD distance_from_x1 = pool.starmap(calculate_ncd, tasks) # 排序与分类逻辑和原代码一致 sorted_idx = np.argsort(np.array(distance_from_x1)) top_k_class = [label for _, label, _ in np.array(training_data)[sorted_idx[:k]]] predict_class = max(set(top_k_class), key=top_k_class.count) return predict_class if __name__ == "__main__": k = 16 training_set = np.array([('TCCCTACACT', 9), ('AGTTGGTATT', 12), ('AGTGGATCAC', 8), ('CGCAAGTGTG', 3), ('CTGTTCCCCC', 7), ('CGACTGGTAA', 10), ('CGCCGAGAAG', 4), ('TAGCTACGAC', 5), ('CCTTGCGCGT', 11), ('CCAAAAGAAA', 12), ('CCTCAGGAGG', 6), ('AGGCCACTTA', 5), ('GCGGGAACGG', 5), ('CTATTACCAA', 2), ('ACACTTTTTT', 8), ('GATGCAGCGT', 1)]) test_set= np.array([('GATGCAGCGT', 3), ('GATGCAGCGT', 0), ('GATGCAGCGT', 3), ('GATGCAGCGT', 4)]) # 提前预计算训练集的压缩长度 training_data = precompute_training_cx(training_set) # 创建全局进程池,复用所有任务 with Pool(cpu_count()) as pool: # 处理所有测试样本(可串行也可并行,核心是复用进程池) results = [process_test_sample(sample, training_data, pool) for sample in test_set] print(results)
关键优化说明
- 全局进程池:在主程序入口通过
with Pool(cpu_count()) as pool创建一次进程池,所有测试样本的内循环计算都复用该池,彻底消除进程创建销毁的开销 - 预计算训练集压缩长度:训练集固定不变,提前计算每个样本的压缩长度,避免每次计算NCD时重复压缩训练样本,节省大量CPU时间
- 减少测试样本重复计算:每个测试样本的压缩长度仅在主进程计算一次,再传递给子进程,避免子进程重复压缩同一个测试样本
- 任务构造优化:将测试样本与每个训练样本的配对构造成独立任务,
starmap自动将元组参数拆分传递给calculate_ncd函数,提升并行效率
内容的提问来源于stack exchange,提问作者DLS
相关产品推荐
相关产品推荐

