【发布时间】:2018-01-02 03:53:06
【问题描述】:
我在 Spark 作业中遇到了一个令我惊讶的错误:
Total size of serialized results of 102 tasks (1029.6 MB) is
bigger than spark.driver.maxResultSize (1024.0 MB)
我的工作是这样的:
def add(a,b): return a+b
sums = rdd.mapPartitions(func).reduce(add)
rdd 有约 500 个分区,func 获取该分区中的行并返回一个大数组(一个 1.3M 双精度数或约 10Mb 的 numpy 数组)。 我想总结所有这些结果并返回它们的总和。
Spark 似乎将 mapPartitions(func) 的总结果保存在内存中(大约 5gb),而不是增量处理它,这只需要大约 30Mb。
除了增加 spark.driver.maxResultSize,有没有办法以增量方式执行 reduce?
更新:实际上,我有点惊讶,有两个结果被保存在内存中。
【问题讨论】:
标签: apache-spark mapreduce rdd