【发布时间】:2016-07-03 13:31:41
【问题描述】:
我加载一个数据集
val data = sc.textFile("/home/kybe/Documents/datasets/img.csv",defp)
我想为此数据建立一个索引
val nb = data.count.toInt
val tozip = sc.parallelize(1 to nb).repartition(data.getNumPartitions)
val res = tozip.zip(data)
不幸的是我有以下错误
Can only zip RDDs with same number of elements in each partition
如果可能的话,如何按分区修改元素的数量?
【问题讨论】:
-
我是一名初学者,所以可能不是最好的解决方案,但理论上可以在两个 RDD 上
zipWithIndex然后执行join(或左连接或右连接,具体取决于你想如何混合两个RDD)使用元素的索引,另一种方法是计算长度差异并用默认值填充空白,但不确定如何使用RDD。 -
我完全忘记了这个伟大的功能!谢谢,但我还是不明白为什么它不起作用。
-
AFAIK 我认为拥有相同数量的分区并不等于拥有两个长度相同的 RDD,我认为
repartition会从文档中对集群中的分区进行洗牌:Reshuffle the data in the RDD randomly to create either more or fewer partitions and balance it across them. This always shuffles all data over the network.,如果我做对了整个事情,那么在你的情况下,如果你需要加入(这会导致另一次改组),那么你的情况就没有多大意义了。
标签: scala apache-spark rdd