【发布时间】:2018-07-06 15:07:16
【问题描述】:
我正在学习 spark,当我在 pyspark shell 中使用以下表达式测试 repartition() 函数时,我观察到一个非常奇怪的结果:所有元素在 repartition() 函数之后落入同一个分区。
在这里,我使用glom() 来了解rdd 内的分区。我期待repartition() 对元素进行洗牌并在分区之间随机分配它们。这只发生在我使用新的分区数
在我的测试过程中,如果我设置新的分区数 > 原始分区,也没有观察到洗牌。我在这里做错了吗?
In [1]: sc.parallelize(range(20), 8).glom().collect()
Out[1]:
[[0, 1],
[2, 3],
[4, 5],
[6, 7, 8, 9],
[10, 11],
[12, 13],
[14, 15],
[16, 17, 18, 19]]
In [2]: sc.parallelize(range(20), 8).repartition(8).glom().collect()
Out[2]:
[[],
[],
[],
[],
[],
[],
[2, 3, 6, 7, 8, 9, 14, 15, 16, 17, 18, 19, 0, 1, 12, 13, 4, 5, 10, 11],
[]]
In [3]: sc.parallelize(range(20), 8).repartition(10).glom().collect()
Out[3]:
[[],
[0, 1],
[14, 15],
[10, 11],
[],
[6, 7, 8, 9],
[2, 3],
[16, 17, 18, 19],
[12, 13],
[4, 5]]
我使用的是 spark 版本 2.1.1。
【问题讨论】:
标签: apache-spark pyspark