如何在Ray分布式集群中使用支持pandas输入的Prophet与Auto ARIMA模型
问题解答
单节点运行判定
是的,如果你直接将原生pandas DataFrame传入未做适配的Prophet、ARIMA这类仅支持单节点运行的模型,相关的训练、预测逻辑只会在触发任务的当前节点运行。原因是这类模型本身没有内置分布式计算的实现,所有运算逻辑都依赖pandas的单节点接口,即便部署在Ray集群环境下,只要没有做任务调度拆分,计算只会占用当前节点的资源,不会自动分发到集群其他节点。
集群适配方案
适配方案按实际使用场景可分为两类:
- 多独立时序并行运算场景
这是绝大多数业务场景的通用情况,比如同时对数千个不同商品、不同门店的独立时序做预测。这种情况不需要修改模型本身的代码,直接用Ray的任务调度能力即可实现分布式运行:- 先用Modin等Ray兼容的分布式数据处理工具加载全量数据,按时序唯一标识拆分成分组,每组数据对应单条独立时序,转换为pandas DataFrame格式
- 用
@ray.remote装饰器封装单模型的训练、预测逻辑,将每个分组的运算任务提交到Ray集群,自动调度到不同节点并行运行 - 最后批量收集所有任务的运行结果即可
示例代码如下:
该方案实现成本极低,资源利用率高,是首选适配方案。import ray import pandas as pd from prophet import Prophet # 连接Ray集群,本地调试时可以不传address参数 ray.init(address="auto") # 封装Prophet运算逻辑为Ray远程任务 @ray.remote def run_single_prophet(train_data: pd.DataFrame, predict_days: int = 30): model = Prophet() model.fit(train_data) future_df = model.make_future_dataframe(periods=predict_days) predict_df = model.predict(future_df) return predict_df # 此处为你拆分好的单时序pandas DataFrame列表 grouped_time_series = [series_1_df, series_2_df, ..., series_n_df] # 批量提交任务到集群并行运行 all_predict_result = ray.get([run_single_prophet.remote(series) for series in grouped_time_series]) - 单条超大时序运算场景
如果你要处理的是单条数据量极大的时序,单机内存/运算速度无法满足需求,可以用以下方式适配:- 超参数调优/交叉验证场景:将不同参数组合的验证任务、不同折的训练任务分发到集群各节点并行运行,筛选最优参数后再完成最终训练
- 长时序训练场景:可以将长时序拆分为多段重叠的短时序,分别在不同节点训练,最后通过加权合并的方式得到最终预测结果
注:Prophet、ARIMA这类时序预测模型本身的训练逻辑是强时序依赖的,如果单条时序没有可拆分的并行维度,无法实现完全的分布式训练,只能通过上述并行优化的方式提升运行效率。
内容的提问来源于stack exchange,提问作者M.Erkin
相关产品推荐
相关产品推荐

