【问题标题】:Spark & Python: strategy for parallelize/map statsmodels sarimaxSpark 和 Python:并行化/映射 statsmodels sarimax 的策略
【发布时间】:2020-03-20 16:39:10
【问题描述】:

我已经为 sarimax(以及一般的时间序列)网格搜索构建了一个 Python 解决方案。

这是一个python类。

在准备好训练和测试集之后,该类将它们存储为对象属性。

稍后,该类构建一个列表,其中每个项目中包含一组用于 statsmodels sarimax 的参数。

然后,这些项目中的每一项都被传递给类 sarimax 方法,用于拟合模型。每个模型都存储在一个列表中,供以后根据用户选择的评分方法进行选择。

类中构建的 sarimax 方法通过对象属性 (self.df_train) 访问训练集

为了并行训练每组参数,我调用 spark 如下:

spark = SparkSession.builder.getOrCreate()
sca = spark.sparkContext

rdd = sca.parallelize(list_of_parameters)
all_models = rdd.map(self.my_sarimax).collect()

它非常适合从 2016 年开始的每月 ts。 但是,如果我尝试喂它更长的 ts,假设从 2014 年开始,火花工作根本不会开始。 它需要一个永恒的“开始”,然后它就会失败。

问题是:

1 - 当我在课堂上运行所有内容时,spark 是否能够理解如何分配此任务?

2 - 集群上的每个节点(worker)能否在需要时轻松找到对象 self.df_train?如果不是,为什么它适用于较短的 ts?我的意思是,这东西很漂亮:训练 9300 多个候选模型平均需要 10 秒。

3 - 如何让它与更长的 ts 一起工作?

【问题讨论】:

    标签: python apache-spark pyspark time-series statsmodels


    【解决方案1】:
    1. spark 能理解如何分配这个任务吗?

      • 是的,虽然每个 spark worker 都在 jvm 下运行,但是如果您将 python 进程分发给 worker(例如在您的情况下为 my_sarimax),每个 worker 将打开单独的 python 进程来运行您的代码。
      • 我没有看到您的完整代码 sn-p,但基于我对问题的理解。您正在准备潜在参数的 rdd,然后将模型和训练数据集广播到分区,然后并行运行所有参数。
    2. 集群上的每个节点(worker)能否在需要时轻松找到对象self.df_train?如果不是,为什么它的工作时间更短?

      • 如果您将类广播到所有分区,则该类将存在于分区/每个工作节点上。
      • 但是如果你广播类,依赖于训练数据,类可能太大,序列化和反序列化需要很多时间,所以编程无法运行。
      • 您的应用程序因 OOM 错误而失败,或者因为传输数据需要很长时间,工作人员没有心跳而被杀死(这可能解释了为什么较小的数据集,您的方法可以正常工作)

    【讨论】:

      猜你喜欢
      • 2017-07-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-12-30
      • 1970-01-01
      • 1970-01-01
      • 2017-07-30
      • 2012-05-13
      相关资源
      最近更新 更多