Python嵌套函数的性能优化与并行化方案咨询
优化方案:先提升串行性能,再实现并行化
一、串行代码优化(基础优化,必须先做)
原代码核心痛点是每次处理df1一行都重复过滤、切片df2,且TzToList存在大量逐行循环操作,先从这两部分优化:
1. 预处理df2,避免重复计算
提前对df2做过滤、分组、索引设置,后续直接调用分组结果,省去重复过滤的开销:
# 预处理df2:过滤NumFZ=0的行,设置ID为索引,按name/group/code分组 df2_preprocessed = df2[df2['NumFZ'] != 0].set_index('ID') # 按联合键分组,生成可直接调用的分组对象 df2_groups = df2_preprocessed.groupby(['name', 'group', 'code'])
2. 重写Filter_df,利用预处理结果
将原isin([x])改为直接相等判断(效率更高),同时直接从分组对象中提取对应数据:
def Filter_df_optimized(row): key = (row['name'], row['group'], row['code']) # 检查分组是否存在,不存在直接返回空列表 if key not in df2_groups.groups: print('No Data at Index:', row.name) return [] df_group = df2_groups.get_group(key) # 按ID范围切片+去重,保留核心列 df_filtered = df_group.loc[row['start']:row['end']].drop_duplicates(subset='ConcFZ', keep='last')[['ConcFZ', 'NumFZ']] if df_filtered.empty: print('No Data at Index:', row.name) return [] return TzToList_optimized(df_filtered)
3. 重写TzToList,用矢量化替代循环
把原函数中的逐行循环改为更高效的批量处理,保留核心逻辑不变:
def TzToList_optimized(df_filtered): # 提取NumFZ=1的ConcFZ,转整数列表 TWTZ = df_filtered[df_filtered['NumFZ'] == 1]['ConcFZ'].astype(int).tolist() # 筛选NumFZ>1的行 df_greater_1 = df_filtered[df_filtered['NumFZ'] > 1] # 处理单行特殊情况 if df_filtered.shape[0] == 1: if df_greater_1.empty: return TWTZ else: tz_list = list(map(int, df_greater_1['ConcFZ'].iloc[0].split(','))) tz_list.sort() TWTZ.append(tz_list[0]) return TWTZ # 批量处理多行NumFZ>1的情况 for _, row in df_greater_1.iterrows(): tz_list = list(map(int, row['ConcFZ'].split(','))) tz_list.sort() # 检查与现有列表无交集则添加最小值 if not set(tz_list) & set(TWTZ): TWTZ.append(tz_list[0]) return TWTZ
二、并行化实现(Win10 32核64线程环境)
Win10下受GIL限制,多线程对CPU密集型任务提升有限,优先选择多进程方案,推荐两种实现:
方案1:concurrent.futures.ProcessPoolExecutor(中小规模数据)
代码改动小,适合数据量未超出内存的场景:
from concurrent.futures import ProcessPoolExecutor # 将df1转为字典列表,方便进程间传递数据 df1_rows = df1.to_dict('records') # 初始化进程池,设置最大工作线程数为核心数(32) with ProcessPoolExecutor(max_workers=32) as executor: # 并行处理每一行数据 results = list(executor.map(Filter_df_optimized, df1_rows)) # 将结果赋值给df1 df1['Filtered'] = results
注意:如果是脚本文件,需将并行代码放在if __name__ == '__main__':块中,避免Win下多进程启动时重复加载代码;Jupyter Lab中可直接运行。
方案2:Dask DataFrame(超大规模数据)
适合数据量超出内存的场景,自动分块并行处理:
步骤1:安装Dask
pip install dask[complete]
步骤2:实现并行逻辑
import dask.dataframe as dd # 将df1转为Dask DataFrame,分块数设为CPU核心数 ddf1 = dd.from_pandas(df1, npartitions=32) # 定义并行处理函数 def dask_filter(row): return Filter_df_optimized(row) # 执行并行计算,指定返回类型为列表 df1['Filtered'] = ddf1.apply(dask_filter, axis=1, meta=('Filtered', object)).compute(scheduler='processes')
Dask的优势是支持超大规模数据的分块处理,避免内存溢出,适合TB级数据集。
三、性能测试建议
- 先验证优化后的串行代码,对比原代码耗时(通常能提升5-10倍)
- 根据数据规模选择并行方案:中小规模用
ProcessPoolExecutor,超大规模用Dask - Win下避免使用多线程方案(如
threading),GIL会限制CPU密集型任务的多线程效率
内容的提问来源于stack exchange,提问作者Rory
相关产品推荐
相关产品推荐

