多进程顺序处理含共享值任务:DataFrame嵌套循环性能优化
问题描述
我有一个遍历pandas DataFrame索引的函数,外层循环遍历i至n(n=200k),内层循环遍历0至i的所有数据。目前单进程处理200k个索引耗时约36小时,希望将其适配为多进程实现。
简化代码如下:
import pandas as pd import numpy as np import random def processing(series): indexes = series.index[2:series.shape[0]] processed_df = pd.DataFrame(index=indexes, columns=['value', 'max_value']) for index in indexes: new_series = series.loc[:index] squared_diff = new_series.diff().dropna() ** 2 integer_index = new_series.index.get_loc(index) sigma_sq = 1 / (integer_index - 1) * sum(squared_diff) max_value = -np.inf value = None for ind in new_series.index[:-1]: temp_integer_index = new_series.index.get_loc(ind) value = (1 / (sigma_sq * np.sqrt(integer_index - temp_integer_index))) if value > max_value: max_value = value processed_df.loc[index, ['value', 'max_value']] = value, max_value return processed_df if __name__ == '__main__': random_list = [] for i in range(0, 100): n = random.randint(1, 50) random_list.append(n) data = np.array(random_list) values = pd.Series(data) results = processing(values) print(results)
我考虑使用Manager管理共享值max_value,将processed_df替换为列表或字典,后续合并为DataFrame,但不确定如何确保每个处理索引的进程按顺序执行,保证max_value仅来自已计算的索引,是否可采用批次处理?望解答。
解决方案
首先明确:你的代码中每个index的计算完全依赖该index之前的原始数据,但各个index之间的计算是独立的——每个index的max_value是自身范围内的最大值,和其他index的结果无关。因此不需要共享max_value,也不需要严格按顺序执行,每个进程可独立处理一批index,最后合并结果即可。
优化方向
1. 预计算减少重复开销
原代码中每个index都重复计算差分平方和、索引位置映射,这部分可以提前预计算:
- 预先计算累积平方差和,每个index的sigma_sq直接用预计算值,避免重复求和
- 预先建立索引到整数位置的映射,避免频繁调用
get_loc
2. 多进程分批次处理
利用multiprocessing.Pool将任务拆分给多个进程,每个进程处理独立的index任务,最后合并所有结果。
优化后代码示例
import pandas as pd import numpy as np import random from multiprocessing import Pool, cpu_count def precompute_values(series): # 预计算累积平方差和 diff_sq = series.diff().dropna() ** 2 cum_sum_diff_sq = np.cumsum(diff_sq) # 预计算索引到整数位置的映射 idx_to_pos = {idx: i for i, idx in enumerate(series.index)} return cum_sum_diff_sq, idx_to_pos def process_single_index(args): index, cum_sum_diff_sq, idx_to_pos = args pos = idx_to_pos[index] # 直接用预计算的累积和计算sigma_sq sigma_sq = cum_sum_diff_sq[pos-1] / (pos - 1) # 向量化计算所有可能的value,直接取最大值 temp_pos_list = np.arange(pos) values = 1 / (sigma_sq * np.sqrt(pos - temp_pos_list)) max_value = values.max() # 对应原代码中最后一次循环的value value = values[-1] return index, value, max_value def processing_parallel(series): indexes = series.index[2:series.shape[0]] # 预计算共享数据 cum_sum_diff_sq, idx_to_pos = precompute_values(series) # 构造任务列表,只传递必要参数(避免传递大对象) tasks = [(idx, cum_sum_diff_sq, idx_to_pos) for idx in indexes] # 使用CPU核心数作为进程数 with Pool(cpu_count()) as pool: results = pool.map(process_single_index, tasks) # 转换为DataFrame并设置索引 processed_df = pd.DataFrame(results, columns=['index', 'value', 'max_value']).set_index('index') return processed_df if __name__ == '__main__': random_list = [random.randint(1,50) for _ in range(100)] values = pd.Series(np.array(random_list)) results_parallel = processing_parallel(values) print(results_parallel)
关键说明
- 无需共享max_value:每个index的max_value是自身范围内的计算结果,和其他进程的任务无关,完全可以独立计算
- 批次处理的实现:
Pool.map会自动将任务分批次分配给进程,也可使用pool.imap自定义批次大小,进一步优化性能 - 向量化替代循环:内层循环用NumPy向量化实现,比原生Python循环快数倍,结合多进程可大幅压缩耗时
内容的提问来源于stack exchange,提问作者Treasure Hunter
相关产品推荐
相关产品推荐

