【问题标题】:Apache Spark Shuffle Writes but no Shuffle ReadsApache Spark Shuffle 写入但没有 Shuffle 读取
【发布时间】: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


【解决方案1】:

Shuffle 是指多个 Spark 阶段之间的数据交换。当所有数据在传输之前从所有执行器序列化时,随机写入出现在阶段结束时。随机读取发生在从所有执行程序收集数据的阶段开始时。为了通过随机读取/写入获得整个画面,您必须在集群模式下运行。在这种情况下,将触发随机读写

编辑

在您的情况下,您在一个实例中运行一个执行程序,这意味着无需从其他执行程序中引入分区,因此不会进行随机读取。关于随机写入,它通过调用缓存在数据帧的保存操作期间出现。

【讨论】:

  • 问题是为什么它写了几个GB却从不读一个位?
  • 例如 spark 数据集 (RDD) 由分区组成。分区将部分数据称为行、行。任务执行期间的所有这些都在执行者之间分配。某些操作需要在执行程序之间随机分配分区(例如,使用相同的键对行进行分组)。
  • 可能我错过了@FaigB 的观点,但它如何解释没有读取的写入?我也只有一个执行者(单节点)
  • @gapvision 添加了关于为什么没有进行随机读取的解释
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-09-07
  • 2016-06-19
  • 2017-10-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-24
相关资源
最近更新 更多