Spark:如何在所有分区中均匀分布我的记录

pra*_*ora 5 apache-spark

我有一个带有 30 条记录的 RDD(键/值对:键是时间戳,值是 JPEG 字节数组)
并且我正在运行 30 个执行程序。我想将此 RDD 重新分区为 30 个分区,以便每个分区获得一条记录并分配给一个执行程序。

当我使用rdd.repartition(30)它时,将我的 rdd 重新分区为 30 个分区,但有些分区获得 2 条记录,有些获得 1 条记录,有些则没有获得任何记录。

在 Spark 中有什么方法可以将我的记录均匀地分布到所有分区。

kma*_*der 0

partitionBy您可以使用该命令并提供多个分区来强制进行新分区。默认情况下,分区器是基于哈希的,但您可以切换到基于范围的分区以获得更好的分布。如果您确实想强制重新分区,可以使用随机数生成器作为分区函数(在 PySpark 中)。

my_rdd.partitionBy(pCount, partitionFunc = lambda x: np.random.randint(pCount))
Run Code Online (Sandbox Code Playgroud)

然而,这经常会导致效率低下的洗牌(节点之间传输大量数据),但如果您的进程受计算限制,那么它是有意义的。