Pandas多进程实现:按月份拆分DataFrame执行数据处理
多进程按月份处理DataFrame的实现方案
需求说明
通过多进程执行数据处理函数,每个进程对应DataFrame中的一个月份;函数需接收三个DataFrame(df1、df2、df3),处理后返回单个DataFrame,最终合并所有进程结果。
问题1:使用starmap实现的参数列表生成方法
starmap可直接接收包含多参数元组的列表,每个元组的元素会依次传入目标函数。只需为每个月份生成对应过滤后的三个DataFrame元组,组成参数列表即可。
完整代码示例
import pandas as pd from multiprocessing import Pool def manipulation(df1_month, df2_month, df3_month): """ 针对单月份数据执行各类数据处理操作 """ # 此处编写你的数据处理逻辑,例如合并、指标计算等 processed_df = pd.concat([df1_month, df2_month, df3_month], axis=1) return processed_df if __name__ == "__main__": # 假设已加载原始df1、df2、df3,且均包含'Date'列 # 获取所有唯一月份列表(可根据实际需求取多表月份的并集) Months = pd.DatetimeIndex(df1['Date']).month.drop_duplicates().tolist() # 生成starmap所需参数列表:每个元素为(当月df1, 当月df2, 当月df3) params_list = [] for m in Months: df1_month = df1[pd.DatetimeIndex(df1['Date']).month == m] df2_month = df2[pd.DatetimeIndex(df2['Date']).month == m] df3_month = df3[pd.DatetimeIndex(df3['Date']).month == m] params_list.append((df1_month, df2_month, df3_month)) # 创建进程池并执行任务 with Pool() as pool: results = pool.starmap(manipulation, params_list) # 合并所有进程的处理结果 final_df = pd.concat(results, ignore_index=True)
问题2:循环创建Process进程的参数传递方法
直接使用multiprocessing.Process创建进程时,可通过args参数传递元组形式的多参数;由于Process不直接返回结果,需借助队列收集每个进程的输出。
完整代码示例
import pandas as pd from multiprocessing import Process, Queue def manipulation(df1_month, df2_month, df3_month, result_queue): """ 针对单月份数据执行各类数据处理操作,结果存入队列 """ # 此处编写你的数据处理逻辑 processed_df = pd.concat([df1_month, df2_month, df3_month], axis=1) result_queue.put(processed_df) if __name__ == "__main__": # 假设已加载原始df1、df2、df3,且均包含'Date'列 Months = pd.DatetimeIndex(df1['Date']).month.drop_duplicates().tolist() # 创建队列用于收集进程结果 result_queue = Queue() processes = [] # 循环创建并启动进程 for m in Months: df1_month = df1[pd.DatetimeIndex(df1['Date']).month == m] df2_month = df2[pd.DatetimeIndex(df2['Date']).month == m] df3_month = df3[pd.DatetimeIndex(df3['Date']).month == m] # 通过args传递参数与结果队列 p = Process(target=manipulation, args=(df1_month, df2_month, df3_month, result_queue)) processes.append(p) p.start() # 等待所有进程执行完成 for p in processes: p.join() # 从队列取出所有结果并合并 results = [] while not result_queue.empty(): results.append(result_queue.get()) final_df = pd.concat(results, ignore_index=True)
注意事项
- 多进程代码必须放在
if __name__ == "__main__":代码块中,避免Windows系统下的进程启动异常。 - 修正了原代码片段中的笔误:原
df2 = df1[pd.DatetimeIndex(df2['Date']).month == m]应改为df2 = df2[pd.DatetimeIndex(df2['Date']).month == m]。 - 若DataFrame数据量极大,建议先按月份拆分数据再传入进程,避免内存占用过高。
内容的提问来源于stack exchange,提问作者Daneel Ank
相关产品推荐
相关产品推荐

