【问题标题】:How to avoid large intermediate result before reduce?在减少之前如何避免大的中间结果?
【发布时间】: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


    【解决方案1】:

    使用reduce 时,Spark 会对驱动程序应用最终缩减。如果func 返回单个对象,这实际上等效于:

    reduce(add, rdd.collect())
    

    您可以使用treeReduce:

    import math
    
    # Keep maximum possible depth
    rdd.treeReduce(add, depth=math.log2(rdd.getNumPartitions()))
    

    toLocalIterator:

    sum(rdd.toLocalIterator())
    

    前者会以增加网络交换为代价递归地合并worker上的分区。您可以使用depth 参数调整性能。

    后者当时只会收集一个分区,但可能需要重新评估rdd,并且大部分工作将由驱动程序执行。

    根据func 中使用的确切逻辑,您还可以通过将矩阵分成块并逐块执行加法来改善工作分配,例如使用BlockMatrices

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-03-16
      • 2017-09-02
      • 2015-08-04
      • 2022-01-04
      • 1970-01-01
      • 2016-03-06
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多