【问题标题】:Spark MLlib - trainImplicit warningSpark MLlib - trainImplicit 警告
【发布时间】:2015-06-30 12:59:44
【问题描述】:

我在使用 trainImplicit 时不断看到这些警告:

WARN TaskSetManager: Stage 246 contains a task of very large size (208 KB).
The maximum recommended task size is 100 KB.

然后任务大小开始增加。我尝试在输入 RDD 上调用 repartition,但警告是相同的。

所有这些警告都来自 ALS 迭代、flatMap 和聚合,例如 flatMap 显示这些警告的阶段的来源(使用 Spark 1.3.0,但它们也在 Spark 1.3.1 中显示):

org.apache.spark.rdd.RDD.flatMap(RDD.scala:296)
org.apache.spark.ml.recommendation.ALS$.org$apache$spark$ml$recommendation$ALS$$computeFactors(ALS.scala:1065)
org.apache.spark.ml.recommendation.ALS$$anonfun$train$3.apply(ALS.scala:530)
org.apache.spark.ml.recommendation.ALS$$anonfun$train$3.apply(ALS.scala:527)
scala.collection.immutable.Range.foreach(Range.scala:141)
org.apache.spark.ml.recommendation.ALS$.train(ALS.scala:527)
org.apache.spark.mllib.recommendation.ALS.run(ALS.scala:203)

从聚合:

org.apache.spark.rdd.RDD.aggregate(RDD.scala:968)
org.apache.spark.ml.recommendation.ALS$.computeYtY(ALS.scala:1112)
org.apache.spark.ml.recommendation.ALS$.org$apache$spark$ml$recommendation$ALS$$computeFactors(ALS.scala:1064)
org.apache.spark.ml.recommendation.ALS$$anonfun$train$3.apply(ALS.scala:538)
org.apache.spark.ml.recommendation.ALS$$anonfun$train$3.apply(ALS.scala:527)
scala.collection.immutable.Range.foreach(Range.scala:141)
org.apache.spark.ml.recommendation.ALS$.train(ALS.scala:527)
org.apache.spark.mllib.recommendation.ALS.run(ALS.scala:203)

【问题讨论】:

  • 能否提供数据和代码示例?
  • 我很惊讶现代框架认为 208KB 是“大”的。想知道这样做的理由是什么......
  • 这是任务的大小,而不是数据的大小。
  • 很可能您的数据出现了偏差,这给一项任务带来了更多负担
  • 看来这些问题,至少对于隐式反馈训练来说,可以放心地忽略。

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


【解决方案1】:

Apache Spark 邮件列表中描述了类似的问题 - http://apache-spark-user-list.1001560.n3.nabble.com/Large-Task-Size-td9539.html

我认为您可以尝试使用分区数量(使用 repartition() 方法),具体取决于您拥有多少主机、RAM、CPU。

还可以尝试通过 Web UI 调查所有步骤,您可以在其中查看阶段数、每个阶段的内存使用情况以及数据位置。

或者,除非一切正常且快速运行,否则不要介意这些警告。

此通知在 Spark 中硬编码(scheduler/TaskSetManager.scala

      if (serializedTask.limit > TaskSetManager.TASK_SIZE_TO_WARN_KB * 1024 &&
          !emittedTaskSizeWarning) {
        emittedTaskSizeWarning = true
        logWarning(s"Stage ${task.stageId} contains a task of very large size " +
          s"(${serializedTask.limit / 1024} KB). The maximum recommended task size is " +
          s"${TaskSetManager.TASK_SIZE_TO_WARN_KB} KB.")
      }

.

private[spark] object TaskSetManager {
  // The user will be warned if any stages contain a task that has a serialized size greater than
  // this.
  val TASK_SIZE_TO_WARN_KB = 100
} 

【讨论】:

    猜你喜欢
    • 2018-07-24
    • 2016-03-07
    • 2015-11-28
    • 1970-01-01
    • 1970-01-01
    • 2015-02-23
    • 2017-12-20
    • 2018-02-17
    • 1970-01-01
    相关资源
    最近更新 更多