【发布时间】: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