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

