如何正确使用Python Multiprocessing结合函数实现批量预测?
问题分析与修正方案
你的代码重复计算单个ID的核心原因是:在pool.apply_async中直接执行了prophet.forecast("ID3", df),这会在创建进程前就提前调用该函数,且固定传入了"ID3",导致所有进程都复用这个固定参数的计算结果。
修正后的代码(apply_async版本)
import multiprocessing as mp import pandas as pd import time # 修正函数参数语法错误后的forecast函数 def forecast(id_val, dataframe): # 你的预测逻辑实现 ... return forecast_dataframe if __name__ == '__main__': id_list = df.id.unique()[1:16] t = time.perf_counter() print("CPU COUNT:", mp.cpu_count()) # 创建进程池,保留4个CPU核心不占用 pool = mp.Pool(mp.cpu_count() - 4) # 正确传递函数对象与动态参数:第一个参数是函数名(不要加括号执行),args传入每个进程的专属参数 results = [pool.apply_async(forecast, args=(idx, df)) for idx in id_list] pool.close() pool.join() print('Extracted in', time.perf_counter() - t) # 初始化结果DataFrame,避免未定义报错 fb_forecast_output = pd.DataFrame() for res in results: fb_forecast_output = pd.concat([fb_forecast_output, res.get()], axis=0) fb_forecast_output.to_csv('./data/output/fb_forecast_output.csv', index=False)
更简洁的pool.map版本
如果希望代码更简洁,可以用functools.partial绑定固定参数,配合pool.map批量处理:
import multiprocessing as mp import pandas as pd import time from functools import partial def forecast(id_val, dataframe): ... return forecast_dataframe if __name__ == '__main__': id_list = df.id.unique()[1:16] t = time.perf_counter() print("CPU COUNT:", mp.cpu_count()) pool = mp.Pool(mp.cpu_count() - 4) # 绑定dataframe参数,让函数变为仅接收id_val的单参数函数 forecast_with_df = partial(forecast, dataframe=df) # map自动遍历id_list,将每个元素传入函数并收集结果 results = pool.map(forecast_with_df, id_list) pool.close() pool.join() print('Extracted in', time.perf_counter() - t) fb_forecast_output = pd.concat(results, axis=0) fb_forecast_output.to_csv('./data/output/fb_forecast_output.csv', index=False)
额外注意事项
- 原始
forecast函数定义存在语法错误:参数不能用字符串"id",应改为合法变量名(比如id_val,避免与Python内置关键字id冲突)。 - 如果数据集
df体积较大,进程间传递数据会产生额外开销,建议改为在每个进程内读取数据文件,减少内存复制成本。
内容的提问来源于stack exchange,提问作者Maxl Gemeinderat
相关产品推荐
相关产品推荐

