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
相关产品推荐
相关产品推荐

