【问题标题】:Running regression on several columns in parallel在多个列上并行运行回归
【发布时间】:2018-05-19 13:27:05
【问题描述】:

我有一个非常宽的带有标签列的数据框。我想独立地为每一列运行逻辑回归。我正在尝试找到并行运行的最有效方法。

+----------+--------+--------+--------+-----+------------+
| features | label1 | label2 | label3 | ... | label30000 |
+----------+--------+--------+--------+-----+------------+

我最初的想法是使用ThreadPoolExecutor,获取每一列的结果,然后加入:

extract_prob = udf(lambda x: float(x[1]), FloatType())

def lr_for_column(argm):
    col_name = argm[0]
    test_res = argm[1]
    lr = LogisticRegression(featuresCol="features", labelCol=col_name, regParam=0.1)
    lrModel = lr.fit(tfidf)
    res = lrModel.transform(test_tfidf)
    test_res = test_res.join(res.select('id', 'probability'), on="id")
    test_res = test_res.withColumn(col_name, extract_prob('probability')).drop("probability")
    return test_res.select('id', col_name)


with futures.ThreadPoolExecutor(max_workers=100) as executor:
    future_results = [executor.submit(lr_for_column, [colname, test_res]) for colname in list_of_label_columns]
    futures.wait(future_results)
    for future in future_results:
       test_res = test_res.join(future.result(), on="id")

但是这种方法的性能不是很好。有没有更快的方法来做到这一点?

【问题讨论】:

  • @user9613318 大约 500000 行,接近 300000 个特征。
  • 数据有多少个分区,你总共分配了多少个核心?还有多少内存/核心?
  • @user9613318 200 个分区,8 节点集群,每个节点有 4 核和 28 GB RAM

标签: python apache-spark pyspark apache-spark-ml


【解决方案1】:

考虑到可用资源,通过使用ThreadPoolExecutor - having 32 cores in total 和 200 个分区,您只能同时处理约 16% 的数据,而这部分只会变得更糟,如果数据成长。

如果您想训练 30000 个模型并使用默认迭代次数(100,在实践中可能很低),您的 Spark 程序将提交大约 3000000 个作业(每次迭代创建一个单独的作业),并且每个作业只提交一小部分可以同时处理 - 这不会给改进带来太大希望,除非您添加更多资源。

尽管有些事情你可以尝试:

  • 确保不必重新计算最终特征。如有必要,将数据写入持久存储并重新加载,并确保将传递给模型的数据缓存。
  • 考虑应用一些降维算法。特征数为 300000 不仅高,而且接近记录数(500000)。它不仅计算量大,而且可能导致严重的过拟合。
  • 如果您决定减少维度,请考虑采样以进一步减少训练数据的大小,从而减少分区数量并提高整体吞吐量。

    如果您的数据中有很强的线性趋势,即使在较小的样本上也应该可见,而不会显着降低精度。

  • 考虑将昂贵的 pyspark.ml 算法替换为不需要多个作业的变体,例如使用来自 spark-sklearn 的某些工具组合(您可以通过在每个模型上拟合 sklearn 模型来创建集成模型分区)。

  • 超额订阅核心。例如,如果您有 4 个物理内核/节点,则允许 8 或 16 个来考虑 IO 等待时间。

【讨论】:

  • 你是如何获得 16% 的?
  • 4 个核心乘以 8 个节点,总共有 32 个核心。如果数据有 200 个分区 -> 200 * 0.16 = 32(每个核心一次可以处理一个分区)。即使包含 IO 等待,它看起来也不好。
  • 这是否意味着我也应该减少分区数?
  • 到大分区通常不好所以可能没有。更广泛的讨论How to calculate the best numberOfPartitions for coalesce?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-17
相关资源
最近更新 更多