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

Pandas大规模多DataFrame join合并性能优化方案咨询

DataFrame合并性能优化问题

场景说明

现有约480个DataFrame,每个DataFrame遵循以下组织规则:

  • 单个DataFrame行数约为600*512
  • 每个DataFrame包含6个相关列,其中timestamp列值唯一,可作为索引使用
  • 各DataFrame可能缺失部分行(即部分时间戳对应的数据),无法直接拼接
  • 执行合并操作前,需要先按索引拼接两个DataFrame,扩展时间戳覆盖范围

初始实现代码

当前采用增量式构建DataFrame的方案,实现代码如下:

# Function to concat index wise the dataframes
def get_joint_dataframe(it, trial_path):
    train_path = os.path.join(trial_path, "train_results_iteration_{0}.csv".format(it)) # 1st dataframe
    test_path = os.path.join(trial_path, "results_iteration_{0}.csv".format(it)) # 2st dataframe

    train_df = pd.read_csv(train_path)
    test_df = pd.read_csv(test_path)

    train_df = pd.concat([train_df, test_df], axis=0, ignore_index=True)# safe to do as I know index sets are disjoint
    train_df = train_df.sort_values(by = 'timestamp').reset_index(drop=True)

    return train_df

# Build incremental dataframe
def get_action_df_speedup(trials):

    df_action_all = None # Incremental dataframe
    policy_count = 0

    for trial in trials:
        for it in range(1, 11):
            df_action = get_joint_dataframe(it, trial)[['timestamp', 'day', 'action']] # In the real case there are 3 additional columns
            df_action = df_action.rename(columns={'action':'policy_{0}'.format(policy_count)})

            columns_to_join = ['policy_{0}'.format(policy_count)] # I do not need to join 'day', just use the one of the first dataframe

            df_action.set_index('timestamp', inplace=True)

            if df_action_all is None:
                df_action_all = df_action
            else:
                df_action_all = df_action_all.join(df_action[columns_to_join], how="inner")
            
            policy_count += 1
        
    df_action_all.reset_index(inplace=True, drop=False)   
    return df_action_all

选择使用join是因为相关技术帖提到,当使用列作为索引时,join的速度通常快于merge。但这段代码运行耗时过长:处理20个DataFrame约需40秒,推算全量数据集处理需要2个半小时,现寻求该场景下的运行速度提升方案。

补充说明

CSV数据样例

timestamp   day   action    reward      Q          Q1         Q2          Q0
0   2.017010e+12    2017/01/03  0.0 0.0 0.008301    0.008301    -1.111009   -0.172822
1   2.017010e+12    2017/01/03  0.0 0.0 0.000000    0.000000    0.000000    0.000000
2   2.017010e+12    2017/01/03  0.0 0.0 0.000000    0.000000    0.000000    0.000000
3   2.017010e+12    2017/01/03  0.0 0.0 0.000000    0.000000    0.000000    0.000000
4   2.017010e+12    2017/01/03  0.0 0.0 0.000000    0.000000    0.000000    0.000000

注意:timestamp列值实际是唯一的,仅展示层面看似重复。此前曾考虑通过列表收集所有DataFrame后统一拼接,但该方案需要重命名所有列,且原则上无法直接拼接所有DataFrame。

自行优化的版本

自行开发了一版理论上耗时更低的解决方案(约120秒,性能提升明显),代码如下:

# Get action dataframe, along with a translator to get the path of a trial
def get_dfs_speedup(trials):

    # Logic
    def manage_df(iteration, trial, policy_count, translator, columns_Q, columns_pol):
        action_df = get_joint_dataframe(iteration, trial)[['timestamp', 'day', 'action', 'Q0', 'Q1', 'Q2']]
        translator['policy_{0}'.format(policy_count[0])] = (trial, iteration)
        dict_transl_action = {'action':'policy_{0}'.format(policy_count[0])}
        dict_transl_Q = {'Q{0}'.format(i):'policy_{1}_Q{0}'.format(i, policy_count[0]) for i in range(3)}
        action_df = action_df.rename(columns={'action':'policy_{0}'.format(policy_count[0])})
        action_df = action_df.rename(columns={'Q{0}'.format(i):'policy_{1}_Q{0}'.format(i, policy_count[0]) for i in range(3)})
        action_df.set_index('timestamp', inplace=True)
        # Update translator
        # Collect indexes which will be used to index dataframe
        columns_Q.extend(list(dict_transl_Q.values()))
        columns_pol.extend(list(dict_transl_action.values()))
        # We now add the df to list
        columns_to_join = list(dict_transl_action.values())
        columns_to_join.extend(list(dict_transl_Q.values()))
        if policy_count[0] == 0:
            columns_to_join.append('day')

        policy_count[0] += 1

        print("Done policy: {0}".format(policy_count[0]))
        return action_df[columns_to_join]
    
    translator = {}
    columns_Q = ['day']
    columns_pol = ['day']
    policy_count = [0]

    all_dfs = [manage_df(iteration, trial, policy_count, translator, columns_Q, columns_pol) for iteration, trial in it.product(list(range(1, 11)), trials)]
    
    # Now we simply concatenate columns wise
    all_dfs = pd.concat(all_dfs, axis=1, ignore_index=False, join='inner')
    # We want both action and q dfs, with translator
    #print(columns_Q, all_dfs.columns)
    Q_df = all_dfs[columns_Q]
    actions_df = all_dfs[columns_pol]
    # We reset_index to get back our timestamp
    Q_df.reset_index(drop=False, inplace=True)
    actions_df.reset_index(drop=False, inplace=True)
    # End
    return Q_df, actions_df, translator

内容的提问来源于stack exchange,提问作者King Powa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 04:51:19