【问题标题】:Spark ML gradient boosted trees not using all nodesSpark ML 梯度提升树不使用所有节点
【发布时间】:2018-08-17 06:47:01
【问题描述】:

我正在使用 pyspark 中的 Spark ML GBTClassifier 在 AWS EMR 集群上具有约 400k 行和约 9k 列的数据帧上训练二进制分类模型。我正在将此与我当前的解决方案进行比较,该解决方案在一个巨大的 EC2 上运行 XGBoost,可以将整个数据帧放入内存中。

我希望我可以在 Spark 中更快地训练(并获得新的观察结果),因为它是分布式/并行的。然而,当观察我的集群时(通过神经节),我看到只有 3-4 个节点有活动的 CPU,而其余的节点只是坐在那里。事实上,从它的外观来看,它可能只使用了一个节点进行实际训练。

我似乎在文档中找不到任何关于节点限制或分区的内容,或者任何似乎与为什么会发生这种情况相关的内容。也许我只是误解了算法的实现,但我认为它的实现方式可以并行化训练以利用 Spark 的 EMR/集群方面。如果不是,那么这样做与仅在单个 EC2 上的内存中进行比较有什么优势吗?我想您不必将数据加载到内存中,但这并不是什么优势。

这是我的代码的一些样板。感谢您的任何想法!

import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType
from pyspark.ml.classification import GBTClassifier
from pyspark.ml.evaluation import BinaryClassificationEvaluator

# Start Spark context:
sc = pyspark.SparkContext()
sqlContext = SparkSession.builder.enableHiveSupport().getOrCreate()

# load data
df = sqlContext.sql('SELECT label, features FROM full_table WHERE train = 1')
df.cache()
print("training data loaded: {} rows".format(df.count()))

test_df = sqlContext.sql('SELECT label, features FROM full_table WHERE train = 0')
test_df.cache()
print("test data loaded: {} rows".format(test_df.count()))


#Create evaluator
evaluator = BinaryClassificationEvaluator()
evaluator.setRawPredictionCol('prob')
evaluator.setLabelCol('label')

# train model
gbt = GBTClassifier(maxIter=100, 
                    maxDepth=3, 
                    stepSize=0.1,
                    labelCol="label", 
                    seed=42)

model = gbt.fit(df)


# get predictions
gbt_preds = model.transform(test_df)
gbt_preds.show(10)


# evaluate predictions
getprob=udf(lambda v:float(v[1]),DoubleType())
preds = gbt_preds.withColumn('prob', getprob('probability'))\
        .drop('features', 'rawPrediction', 'probability', 'prediction')
preds.show(10)

auc = evaluator.evaluate(preds)
auc

旁注:我使用的表格已经矢量化。该模型使用此代码运行,它运行缓慢(大约 10-15 分钟的训练时间)并且只使用 3-4 个(或者可能只使用其中一个)内核。

【问题讨论】:

  • 所以您使用的是 EMR,对吗?您是否使用spark-submit 提交作业?如果是这样,您能否发布用于提交作业的整个命令?
  • 是的,我只是使用spark-submit spark_gbt.py(其中spark_gbt.py 是上面脚本的名称)从命令行ssh'd 到主节点。您想知道集群上的任何特定设置吗?我没有自己设置这些,但如果你认为这相关,我当然可以追踪它们。我会说我在同一个集群上运行的其他 pyspark 作业使用所有节点。
  • 我想弄清楚您是使用本地模式还是实际使用分布式模式。据我所知,如果您没有指定--master yarn,那么它将在 EMR 上以本地模式运行。我还注意到另一件事 - 您没有在代码中缓存任何数据。在运行算法之前,您应该运行 df.cache() 甚至可能是 df.count() 来初始化缓存。
  • 它正在运行分布式模式。我用--master yarn 跑来仔细检查,行为是一样的。我还添加了df.cache()df.count() 行,没有任何改变。 (请参阅上面的编辑位置。)
  • 我想,更清楚地说,我的问题是:其他人在 Spark ML 中训练模型时会看到同样的行为吗?

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


【解决方案1】:

感谢上述 cmets 的澄清。

Spark 的实现不一定比 XGBoost 快。事实上,我会期待你所看到的。

最大的因素是 XGBoost 在设计和编写时特别考虑了梯度提升树。另一方面,Spark 的用途更广泛,而且很可能没有 XGBoost 所具有的那种优化。请参阅here,了解 XGBoost 和 scikit-learn 的分类器算法实现之间的区别。如果您想真正了解细节,可以阅读论文,甚至可以阅读 XGBoost 和 Spark 实现背后的代码。

请记住,XGBoost 也是并行/分布式的。它只是在同一台机器上使用多个线程。当数据不适合单台机器时,Spark 可帮助您运行算法。

我能想到的其他几个小问题是:a) Spark 确实有一个重要的启动时间。不同机器之间的通信也可以加起来。 b) XGBoost 是用 C++ 编写的,通常非常适合数值计算。

至于为什么 Spark 只使用 3-4 个核心,这取决于你的数据集大小,它是如何跨节点分布的,spark 启动的执行器数量是多少,占用了哪个阶段大多数时间,内存配置等。您可以使用 Spark UI 尝试找出发生了什么。如果不查看数据集,很难说出为什么会发生这种情况。

希望对您有所帮助。

编辑:我刚刚找到了一个比较简单 Spark 应用程序与独立 Java 应用程序之间的执行时间的好答案 - https://stackoverflow.com/a/49241051/5509005。同样的原则也适用于此,事实上,由于 XGBoost 已高度优化,因此更是如此。

【讨论】:

  • 这很有帮助。谢谢你。我隐约想到了这些权衡,但得到确认是件好事。我的问题仍然是:其他人看到同样的事情吗?我在使用LogisticRegression 分类器时也看到了它。就像你说的,XGBoost 是跨核心并行化的,所以你会认为 Spark 至少会尝试跨节点并行化来竞争。我再次尝试了 1400 万行,花了一个多小时,但仍然只使用一个节点(不是主节点)进行训练。
  • 在训练 ML 模型时,我必须检查 Spark UI 的唯一一次是运行 Gradient Boosted Trees 模型 - 与您在原始帖子中看到的相同。我记得有许多节点在工作——它肯定多于 1 并且少于集群中的节点总数。它并没有打扰我,因为集群非常大——25-30 个节点。
  • 更新:我仍然无法让它使用超过 1 个节点进行训练,但如果我多线程处理几个不同的 spark-submit 命令,那么它将向每个节点发送一个模型训练作业。这只是有用的,因为我实际上每次都在尝试训练大约 20 个这些模型(有点像迷你网格搜索)。我仍然很好奇是否有人知道 Spark ML 培训是否应该是分布式的,但据我所知不是。感谢所有的帮助。
猜你喜欢
  • 1970-01-01
  • 2018-09-10
  • 1970-01-01
  • 2020-12-21
  • 2013-12-08
  • 2019-06-27
  • 2018-08-21
  • 2015-07-12
  • 1970-01-01
相关资源
最近更新 更多