【发布时间】:2017-03-13 22:49:10
【问题描述】:
我有一个执行大型连接的 Spark 应用程序
val joined = uniqueDates.join(df, $"start_date" <= $"date" && $"date" <= $"end_date")
然后将生成的 DataFrame 聚合为一个可能有 13k 行的数据帧。在加入过程中,作业失败并显示以下错误消息:
Caused by: org.apache.spark.SparkException: Job aborted due to stage failure: Total size of serialized results of 78021 tasks is bigger than spark.driver.maxResultSize (2.0 GB)
这发生在之前没有设置spark.driver.maxResultSize,所以我设置了spark.driver.maxResultSize=2G。然后,我对连接条件稍作更改,错误再次出现。
编辑:在调整集群大小时,我还将 .coalesce(256) 中的 DataFrame 假定的分区数量翻了一番,达到了 .coalesce(512),所以我不能确定不是因为这个.
我的问题是,既然我没有向司机收集任何东西,为什么spark.driver.maxResultSize 在这里很重要?驱动程序的内存是否用于我不知道的连接中的某些内容?
【问题讨论】:
-
遇到同样的问题,请问您有什么进展吗?
-
@user4601931 你能粘贴你正在运行的实际 scala 代码吗?
val joined = uniqueDates.join(df, $"start_date" <= $"date" && $"date" <= $"end_date")行不会运行任何作业。您必须进行一些触发工作的转换。 -
您能检查一下
joined中有多少个分区吗?像joined.queryExecution.toRdd.getNumPartitions这样的东西。我很好奇你为什么有78021 tasks。是否更好的解决方案是减少连接中数据集的分区数量? -
@JacekLaskowski 不幸的是,我已经没有这个项目的代码了,而且时间太长了,我已经忘记了它的大部分内容。抱歉,但感谢您对这个问题再次感兴趣。
-
@JacekLaskowski 我无法在此处显示查询计划,但它崩溃的阶段包含 +3000 个任务。很多
FileScanRDD后跟MapPartitionsRDD。然后很多UnionRDD。最后对所有联合的结果进行不同的操作。但是没有(广播)加入或收集......我当然可以看到为什么这个执行计划不理想,但不是spark.driver.maxResultSize进来的地方。当--deploy-mode cluster设置时没有崩溃。
标签: scala apache-spark memory apache-spark-sql