【问题标题】:How to turn a known structured RDD to Vector如何将已知的结构化 RDD 转换为 Vector
【发布时间】: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


【解决方案1】:

一个关键的事情:你真的需要 Vector 吗?地图可能更合适。

  • 如果你真的需要局部Vector,你首先需要使用.collect(),然后将局部转换成Vector。当然,你应该有足够的内存。但这里真正的问题是在哪里可以找到可以从 (index, value) 对有效构建的 Vector。据我所知,Spark MLLib 有自己的类org.apache.spark.mllib.linalg.Vectors,它可以从索引和值数组甚至元组创建Vector。在引擎盖下它使用breeze.linalg。所以这可能是你最好的开始。

  • 如果您只需要订购,您可以使用.orderByKey(),因为您已经拥有RDD[(K,V)]。这样你就订购了流。这并不严格遵循您的意图,但也许它可能更适合。现在您可以通过 .reduceByKey() 删除具有相同键的元素,只生成结果元素。

  • 1234563 ' 缺少索引的元素。

希望这已经足够清楚了。

【讨论】:

  • 这些似乎都比我回答中的解决方案复杂。特别是 sort + reduce 或 sort + flatmap 效率非常低,而 collectAsMap 避免了需要足够内存的问题(当然,所有解决方案都需要假设最终的 Vector 将适合内存)
  • 好的,这可能是因为我的具体情况——我通常使用大型 RDD,所以你的解决方案对我来说看起来很糟糕,因为 .updated()。我认为@Noam 必须重新考虑他是否真的需要 Vector。
  • 嗯...对我来说“适合记忆”通常是错误的建议 ;-)。
  • 折叠可以做任何事情(这实际上是真的:cs.nott.ac.uk/~gmh/fold.pdf
  • 在我的情况下,我知道我的 RDD 很小,因为它已经被键减少了我想要一个向量来更新一个累加器而不是更新多个累加器
【解决方案2】:
rdd.collectAsMap.foldLeft(Vector[Int]()){case (acc, (k,v)) => acc updated (k, v)}

将 RDD 变成 Map。然后对其进行迭代,构建一个 Vector。

您可以使用 justt collect(),但如果有许多重复的具有相同键的元组可能不适合内存。

【讨论】:

  • >8-[] @Paul 在这里使用.updated() 是否可以接受它实际上提供了副本?我了解people want keeping things simple 但是...
  • 向量是不可变的。真的没有别的选择了。此外,Vector 是一种持久数据结构,它实际上并不复制所有内容(新的“副本”与旧 Vector 共享其大部分结构)
  • 您的答案最接近我所寻找的,并为我指明了解决方案
  • 这次讨论对我来说最好的部分是关于折叠方法的好点,谢谢@Paul。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-11-30
  • 1970-01-01
  • 2015-10-28
  • 2017-06-02
  • 1970-01-01
  • 2015-03-25
  • 2020-03-11
相关资源
最近更新 更多