我有一个带有 30 条记录的 RDD(键/值对:键是时间戳,值是 JPEG 字节数组)
并且我正在运行 30 个执行程序。我想将此 RDD 重新分区为 30 个分区,以便每个分区获得一条记录并分配给一个执行程序。
当我使用rdd.repartition(30)它时,将我的 rdd 重新分区为 30 个分区,但有些分区获得 2 条记录,有些获得 1 条记录,有些则没有获得任何记录。
在 Spark 中有什么方法可以将我的记录均匀地分布到所有分区。
partitionBy您可以使用该命令并提供多个分区来强制进行新分区。默认情况下,分区器是基于哈希的,但您可以切换到基于范围的分区以获得更好的分布。如果您确实想强制重新分区,可以使用随机数生成器作为分区函数(在 PySpark 中)。
my_rdd.partitionBy(pCount, partitionFunc = lambda x: np.random.randint(pCount))
Run Code Online (Sandbox Code Playgroud)
然而,这经常会导致效率低下的洗牌(节点之间传输大量数据),但如果您的进程受计算限制,那么它是有意义的。
| 归档时间: |
|
| 查看次数: |
6612 次 |
| 最近记录: |