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

DataFrame行多进程处理卡顿问题排查

问题分析与修复

你的multip()函数存在两个核心错误,导致多进程逻辑完全失效甚至陷入异常:

错误1:主进程提前执行分片处理,未真正启用多进程

[iterate_rows(df) for df in listParts]这行代码是在主进程里遍历所有分片并执行iterate_rows,相当于已经在单进程模式下完成了所有分片的拼接工作,完全没有用到后续创建的进程池。这不仅浪费多进程资源,还导致传入pool.map的是已经处理好的DataFrame,而非原始分片。

错误2:pool.map的函数与参数不匹配

row_to_df函数的设计是接收单一行数据(元组),但你现在传入的是已经拼接好的DataFrame分片。类型不匹配会导致row_to_df执行逻辑混乱,进而引发后续拼接失败甚至无限循环。


修正后的multip()函数

def multip():
    Nrows = df1.shape[0]
    Ncores = mp.cpu_count()
    Nparts = min([Nrows, Ncores])
    listParts = np.array_split(df1, Nparts)
    pool = mp.Pool(Ncores)
    # 直接让进程池调用iterate_rows处理每个分片,而非主进程提前处理
    processed_dfs = pool.map(iterate_rows, listParts)
    dfY = pd.concat(processed_dfs)
    pool.close()
    pool.join()
    return dfY

额外优化建议

  • 避免在子进程中使用display(dfY),多进程环境下的交互式输出可能引发异常,建议改为打印日志或仅在主进程输出结果。
  • 对于大规模数据,多次pd.concat效率较低,可让子进程返回数据列表而非DataFrame,最后一次性拼接以提升性能:
def row_to_list(row):
    # 返回列表而非DataFrame,减少内存开销
    dfX = df2.copy()
    dfX.loc[:, "Product"] = row[0]
    dfX.loc[:, "Store"] = row[1]
    return dfX.values.tolist()

def iterate_rows_to_list(df):
    result = []
    for row in df.itertuples(index=False):
        result.extend(row_to_list(row))
    return result

def multip_optimized():
    Nrows = df1.shape[0]
    Ncores = mp.cpu_count()
    Nparts = min([Nrows, Ncores])
    listParts = np.array_split(df1, Nparts)
    pool = mp.Pool(Ncores)
    processed_lists = pool.map(iterate_rows_to_list, listParts)
    # 合并所有列表后一次性转为DataFrame
    combined_data = []
    for lst in processed_lists:
        combined_data.extend(lst)
    dfY = pd.DataFrame(combined_data, columns=["Month", "Product", "Store"])
    pool.close()
    pool.join()
    return dfY

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 07:50:48