【问题标题】:Why does DataFrame.stat.approxQuantile fail with size of serialized results of n tasks (1030.8 MB) is bigger than spark.driver.maxResultSize为什么 DataFrame.stat.approxQuantile 失败,n 个任务的序列化结果大小(1030.8 MB)大于 spark.driver.maxResultSize
【发布时间】:2019-06-02 16:42:44
【问题描述】:

val postsQuantiles = posts.stat.approxQuantile("_score", Array(0.25, 0.75), 0.0) 失败并出现以下错误。我显然可以设置spark.driver.maxResultSize 来克服这个错误,但我很好奇为什么这会收集数据给驱动程序?

[Stage 3:==================>                                      (7 + 15) / 22]19/06/01 20:46:30 ERROR TaskSetManager: Total size of serialized results of 18 tasks (1030.8 MB) is bigger than spark.driver.maxResultSize (1024.0 MB)
19/06/01 20:46:30 ERROR TaskSetManager: Total size of serialized results of 19 tasks (1087.7 MB) is bigger than spark.driver.maxResultSize (1024.0 MB)
19/06/01 20:46:30 ERROR TaskSetManager: Total size of serialized results of 20 tasks (1145.6 MB) is bigger than spark.driver.maxResultSize (1024.0 MB)
19/06/01 20:46:30 ERROR TaskSetManager: Total size of serialized results of 21 tasks (1203.5 MB) is bigger than spark.driver.maxResultSize (1024.0 MB)
19/06/01 20:46:30 ERROR TaskSetManager: Total size of serialized results of 22 tasks (1261.4 MB) is bigger than spark.driver.maxResultSize (1024.0 MB)
[Stage 3:====================================>                    (14 + 8) / 22]org.apache.spark.SparkException: Job aborted due to stage failure: Total size of serialized results of 18 tasks (1030.8 MB) is bigger than spark.driver.maxResultSize (1024.0 MB)
  at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1599)
  at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1587)
  at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1586)
  at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
  at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
  at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1586)
  at org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:831)
  at org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:831)
  at scala.Option.foreach(Option.scala:257)
  at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:831)
  at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:1820)
  at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1769)
  at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1758)
  at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48)
  at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:642)
  at org.apache.spark.SparkContext.runJob(SparkContext.scala:2034)
  at org.apache.spark.SparkContext.runJob(SparkContext.scala:2131)
  at org.apache.spark.rdd.RDD$$anonfun$fold$1.apply(RDD.scala:1092)
  at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
  at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
  at org.apache.spark.rdd.RDD.withScope(RDD.scala:363)
  at org.apache.spark.rdd.RDD.fold(RDD.scala:1086)
  at org.apache.spark.rdd.RDD$$anonfun$treeAggregate$1.apply(RDD.scala:1155)
  at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
  at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
  at org.apache.spark.rdd.RDD.withScope(RDD.scala:363)
  at org.apache.spark.rdd.RDD.treeAggregate(RDD.scala:1131)
  at org.apache.spark.sql.execution.stat.StatFunctions$.multipleApproxQuantiles(StatFunctions.scala:102)
  at org.apache.spark.sql.DataFrameStatFunctions.approxQuantile(DataFrameStatFunctions.scala:100)
  at org.apache.spark.sql.DataFrameStatFunctions.approxQuantile(DataFrameStatFunctions.scala:75)
  ... 56 elided

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    approxQuantile 方法遵循用于计算近似分位数的 Greenwald-Khanna 算法(基于他们的文档 https://spark.apache.org/docs/latest/api/scala/index.html#org.apache.spark.sql.DataFrameStatFunctions)。它允许您选择相对误差项。

    在他们警告您的文档中,选择相对错误0.0(如您所见)可能非常昂贵,这正是您所看到的。该算法比直接分位数更适合近似分位数。之所以要拉出这么多数据,是因为要计算直接分位数,它至少需要将列中的所有数据拉入驱动程序。

    您可以从已发表的论文中阅读有关该算法的更多信息:http://infolab.stanford.edu/~datar/courses/cs361a/papers/quantiles.pdf

    为了克服这个问题,我建议使用具有适当小值的相对误差项,这会给您“足够接近”的信心。

    【讨论】:

    • 在 Spark 中是否有其他方法可以计算分位数而不需要近似值?
    • 您可以使用 Spark SQL,这很有效。例如,如果您有一个 spark 数据框 df,并且您希望计算列 foo 的中位数(50% 分位数),您可以这样做 df.registerTempTable("df") 和 spark.sql("select percentile(foo, 0.5) from df")
    猜你喜欢
    • 1970-01-01
    • 2018-06-08
    • 1970-01-01
    • 1970-01-01
    • 2016-03-06
    • 2019-04-03
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多