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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 18:05:21