【问题标题】:RDD intersectionsRDD 路口
【发布时间】:2017-04-14 01:58:29
【问题描述】:

我有一个关于两个 RDD 之间交集的查询。

我的第一个 RDD 有一个这样的元素列表:

A = List(1,2,3,4), List(4,5,6), List(8,3,1),List(1,6,8,9,2)

而第二个RDD是这样的:

B = (1,2,3,4,5,6,8,9)

(我可以将 B 作为 Set 存储在内存中,但不是第一个。)

我想将 A 中的每个元素与 B 相交

List(1,2,3,4).intersect((1,2,3,4,5,6,8,9))
List(4,5,6).intersect((1,2,3,4,5,6,8,9))
List(8,3,1).intersect((1,2,3,4,5,6,8,9))
List(1,6,8,9,2).intersect((1,2,3,4,5,6,8,9))

如何在 Scala 中做到这一点?

【问题讨论】:

  • 看看这个:spark.apache.org/docs/1.4.0/api/java/org/apache/spark/rdd/…我花了 10 秒在谷歌上找到。
  • 谢谢!但我怀疑它在scala中是否有效。我收到一条错误消息,说 RDD 不是传递给交点的有效参数
  • 对不起,如果我的问题不清楚:) 但我想做的是一个 RDD 的每个元素与另一个 RDD 的交集。 val output = a4.map(v => v.intersect(a6)) 我得到这个 [错误] /LabWork/BigData/MiniProject/TwitterProject/src/main/scala/community/spark/twitter/TestPrint.scala:41 : 错误:值相交不是 Iterable[String] 的成员
  • 不要与 RDD 相交。你说你可以将它存储在内存中。这样做,将其存储为 Seq,然后 a.map(_.interesect(bAsSeq))

标签: scala intersection rdd


【解决方案1】:

val result = rdd.map( x => x.intersect(B))

Brdd 必须具有相同的类型(在本例中为 List[Int])。另外,请注意,如果 B 很大但适合内存,您可能希望将其广播为 documented here

scala> val B = List(1,2,3,4,5,6,8,9)
B: List[Int] = List(1, 2, 3, 4, 5, 6, 8, 9)

scala> val rdd = sc.parallelize(Seq(List(1,2,3,4), List(4,5,6), List(8,3,1),List(1,6,8,9,2)))
rdd: org.apache.spark.rdd.RDD[List[Int]] = ParallelCollectionRDD[0] at parallelize at <console>:21

scala> rdd.map( x => x.intersect(B)).collect.mkString("\n")
res3: String = 
List(1, 2, 3, 4)
List(4, 5, 6)
List(8, 3, 1)
List(1, 6, 8, 9, 2)

【讨论】:

  • 感谢 KrisP!我正在按照您的建议进行操作,但之前遇到了类型转换错误。现在可以了。
猜你喜欢
  • 1970-01-01
  • 2019-04-25
  • 2016-05-06
  • 1970-01-01
  • 2019-08-11
  • 2020-04-12
  • 2021-10-16
  • 2017-02-17
  • 1970-01-01
相关资源
最近更新 更多