【发布时间】:2016-10-02 05:22:30
【问题描述】:
如果我的 hadoop 块大小为 128 MB,而我的文件为 30 MB。 我运行 spark 的集群是一个 4 节点集群,总共有 64 个内核。
现在我的任务是运行随机森林或梯度提升算法,在此基础上使用参数网格和 3 折交叉验证。
几行代码:
import org.apache.spark.ml.tuning.{ParamGridBuilder, TrainValidationSplit, CrossValidator}
import org.apache.spark.ml.regression.GBTRegressor
val gbt_model = new GBTRegressor().setLabelCol(target_col_name).setFeaturesCol("features").setMaxIter(2).setMaxDepth(2).setMaxBins(1700)
var stages: Array[org.apache.spark.ml.PipelineStage] = index_transformers :+ assembler :+ gbt_model
val paramGrid = new ParamGridBuilder().addGrid(gbt_model.maxIter, Array(100, 200)).addGrid(gbt_model.maxDepth, Array(2, 5, 10)).build()
val cv = new CrossValidator().setEstimator(pipeline).setEvaluator(new RegressionEvaluator).setEstimatorParamMaps(paramGrid).setNumFolds(5)
val cvModel = cv.fit(df_train)
我的文件有
输入: 10 个离散/字符串/字符特征 + 2 个整数特征
输出:一个整数响应/输出变量
这需要 4 个多小时才能在我的集群上运行。我观察到的是,我的代码仅在 1 个节点上运行,只有 3 个容器。
问题:
- 我可以在这里做些什么来确保我的代码在所有四个节点上运行或使用尽可能多的内核进行快速计算。
- 在对数据进行分区(Scala 中的 DataFrame 和 Hadoop 集群上的 csv 文件)方面我可以做些什么来提高速度和计算能力
问候,
【问题讨论】:
标签: hadoop apache-spark apache-spark-mllib apache-spark-ml