为什么pandas DataFrame用multiprocessing做apply并行化一直无响应,如何解决
pandas大数据量apply并行卡住无输出解决方案(无需重构现有处理函数)
- 方案1:修正multiprocessing调用规范(零业务代码改动)
90%以上的并行卡住问题都是未遵循多进程调用规范导致,将你的并行逻辑调整为以下写法即可:import pandas as pd import numpy as np import multiprocessing from your_module import add_features # 你已存为独立文件的处理函数 def parallel_process(df, n_cores=None): if not n_cores: n_cores = multiprocessing.cpu_count() # 拆分数据集 df_splits = np.array_split(df, n_cores) # 新增maxtasksperchild参数避免子进程内存泄漏卡顿 with multiprocessing.Pool(n_cores, maxtasksperchild=1) as pool: processed_splits = pool.map(add_features, df_splits) return pd.concat(processed_splits, ignore_index=True) # 所有多进程逻辑必须放在__main__入口块内,禁止直接写在全局作用域 if __name__ == "__main__": # 加载你的原始数据 raw_df = pd.read_csv("your_news_data.csv") # 执行并行处理 result_df = parallel_process(raw_df) # 保存结果 result_df.to_csv("cleaned_news_data.csv", index=False) - 方案2:使用封装好的并行库(改动量最小)
直接替换原有apply调用为第三方库封装好的并行方法,不需要手动处理分片、进程调度逻辑,完全兼容你现有add_features函数:- 先安装依赖:
pip install pandarallel - 代码调整示例:
import pandas as pd import multiprocessing from pandarallel import pandarallel from your_module import add_features # 初始化并行环境,开启进度条方便观察运行状态 pandarallel.initialize(nb_workers=multiprocessing.cpu_count(), progress_bar=True) raw_df = pd.read_csv("your_news_data.csv") # 仅需将原有的apply改为parallel_apply即可 result_df = raw_df.parallel_apply(add_features, axis=1) result_df.to_csv("cleaned_news_data.csv", index=False) - 先安装依赖:
- 方案3:切换为多线程模式(兼容所有环境)
如果多进程模式依旧存在兼容问题,可以切换为多线程实现,代码调整量极小:import pandas as pd import numpy as np from concurrent.futures import ThreadPoolExecutor from your_module import add_features def parallel_process_thread(df, n_workers=8): df_splits = np.array_split(df, n_workers) with ThreadPoolExecutor(max_workers=n_workers) as executor: processed_splits = list(executor.map(add_features, df_splits)) return pd.concat(processed_splits, ignore_index=True) if __name__ == "__main__": raw_df = pd.read_csv("your_news_data.csv") result_df = parallel_process_thread(raw_df) result_df.to_csv("cleaned_news_data.csv", index=False)
测试时可先抽取1万行以内的子数据集验证逻辑可用性,确认运行正常后再执行全量91万行数据的处理,避免无效等待。
内容的提问来源于stack exchange,提问作者Thomas GF
相关产品推荐
相关产品推荐

