【问题标题】:what is the difference between rdd.repartition() and partition size in sc.parallelize(data, partitions)rdd.repartition() 和 sc.parallelize(data, partitions) 中的分区大小有什么区别
【发布时间】:2015-08-22 02:36:12
【问题描述】:

我正在浏览 spark 的文档。我对 rdd.repartition() 函数和我们在 sc.parallelize() 中的上下文初始化期间传递的分区数有点困惑。

我的机器上有 4 个内核,如果我 sc.parallelize(data, 4) 一切正常,但是当我 rdd.repartition(4) 并应用 rdd.mappartitions(fun) 有时分区没有数据并且我的在这种情况下,函数会失败。

所以,只是想了解这两种分区方式有什么区别。

【问题讨论】:

  • 两者都可能导致空分区。当您编写一个打算与mapPartitions 一起使用的函数时,您应该简单地考虑到这一点。

标签: python apache-spark rdd


【解决方案1】:

通过调用repartition(N),spark 将进行随机播放以更改分区数(默认情况下会生成具有该分区数的 HashPartitioner)。当您使用所需数量的分区调用 sc.parallelize 时,它会将您的数据(或多或少)平均分配到切片之间(实际上类似于范围分区器),您可以在 ParallelCollectionRDD 内部的 slice 函数中看到这一点.

话虽如此,sc.parallelize(data, N)rdd.reparitition(N)(实际上几乎任何形式的数据读取)都可能导致 RDD 的分区为空(这是 @987654327 的一个非常常见的错误来源@ 代码,所以我偏向于 spark-testing-base 中的 RDD 生成器,以创建具有空分区的 RDD)。对于大多数函数来说,一个非常简单的解决方法就是检查你是否传入了一个空迭代器,在这种情况下只返回一个空迭代器。

【讨论】:

    猜你喜欢
    • 2017-12-05
    • 2018-06-18
    • 2020-12-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-11
    • 2010-10-04
    相关资源
    最近更新 更多