【发布时间】:2015-02-17 18:33:08
【问题描述】:
假设我有一个包含 (Int, Int) 元组的 RDD。 我希望把它变成一个向量,其中元组中的第一个 Int 是索引,第二个是值。
任何想法我该怎么做?
我更新了我的问题并添加了我的解决方案以澄清: 我的RDD已经被key减少了,key的数量是已知的。 我想要一个向量来更新单个累加器而不是多个累加器。
我的最终解决方案是:
reducedStream.foreachRDD(rdd => rdd.collect({case (x: Int,y: Int) => {
val v = Array(0,0,0,0)
v(x) = y
accumulator += new Vector(v)
}}))
使用文档中累加器示例中的Vector。
【问题讨论】:
-
如果索引在您的 RDD 中出现两次(或更多)会发生什么?
-
@Paul 我的 RDD 在那里被密钥减少了,因为我知道我不会有重复的密钥,说,在我的情况下,我认为这并不重要,我只会更新累加器几次而不是一次
标签: scala vector apache-spark distributed-computing rdd