【发布时间】:2017-03-01 08:51:17
【问题描述】:
又是一个火花洗牌问题......
我在一个相当大的实例(160GB RAM,40 个核心)上设置了一个节点。我确实想用不同的参数(基于相同的数据框)训练 250 个 ALS 模型。
在调查磁盘 I/O 问题多天后,我发现了这个对话:http://apache-spark-developers-list.1001551.n3.nabble.com/Eliminating-shuffle-write-and-spill-disk-IO-reads-writes-in-Spark-td16955.html
我确实像他们说的那样将我的spark.local.dir 指向了一个 RAM 磁盘。我认为它完成了这项工作,所有核心现在都得到了很好的使用。好!
但是,我确实在 executors 选项卡中看到了一些我无法理解的内容:
看起来只有 Shuffle 写入,没有读取。简单的问题:为什么?有必要吗?如果没有:我该如何避免呢?
执行的主要块:
val modelsAndResults = parameters.par.map( e => {
var auc = 0d
splits.zipWithIndex.foreach { case ((training, validation), splitIndex) =>
val trainingDataset = spark.createDataFrame(training, schema).cache()
val validationDataset = spark.createDataFrame(validation, schema).cache()
val als = baseALS.copy(ParamMap(
baseALS.rank -> e._2,
baseALS.maxIter -> e._1,
baseALS.alpha -> e._4,
baseALS.regParam -> e._5,
baseALS.nonnegative -> e._6));
//println(s"$e evaluating $splitIndex ...")
val localAUC = new BinaryClassificationEvaluator().evaluate(
validationDataset
.join(als.fit(trainingDataset).transform(validationDataset.drop("time")), Seq("user","article"))
.withColumn("label", when($"time" >= e._2, 1d).otherwise(0d)).drop("time")
.withColumn("rawPrediction", $"prediction".cast(DoubleType))
.select("label","rawPrediction"))
auc += localAUC
trainingDataset.unpersist()
validationDataset.unpersist()
//println(s"$e evaluating $splitIndex ... localAUC: $localAUC")
}
val finalAUC = auc / splits.size
val csv = CSVWriter.open("cf_gridsearch_results_12345.csv", append = true)
csv.writeRow(List(e._1,e._2,e._3,e._4,e._5,e._6,finalAUC))
csv.close()
println((e._1,e._2,e._3,e._4,e._5,e._6,finalAUC))
(e._1,e._2,e._3,e._4,e._5,e._6,finalAUC)
})
【问题讨论】:
-
请提供代码 sn-p 发生主要操作的位置(如收集、使用键操作)
标签: apache-spark apache-spark-mllib