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

多进程顺序处理含共享值任务: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 03:48:29