You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 08:55:37