【问题标题】:Extract elements of lists in an RDD提取 RDD 中列表的元素
【发布时间】:2017-03-01 04:48:53
【问题描述】:

我想要达到的目标

我正在使用 Spark 和 Scala。我有两个 Pair RDD。

rdd1 : RDD[(String, List[String])]
rdd2 : RDD[(String, List[String])]

两个 RDD 都以它们的第一个值连接。

val joinedRdd = rdd1.join(rdd2)

所以生成的 RDD 的类型是 RDD[(String, (List[String], List[String]))]。我想映射这个 RDD 并提取两个列表的元素,以便生成的 RDD 只包含两个列表的这些元素。


示例

rdd1 (id, List(a, b))
rdd2 (id, List(d, e, f))
wantedResult (a, b, d, e, f)

天真的方法

我的幼稚方法是直接使用(i) 处理每个元素,如下所示:

val rdd = rdd1.join(rdd2)
    .map({ case (id, lists) => 
        (lists._1(0), lists._1(1), lists._2(0), lists._2(2), lists._2(3)) })

/* results in RDD[(String, String, String, String, String)] */

有没有一种方法可以获取每个列表的元素,而无需单独处理每个元素?类似“lists._1.extractAll”的东西。有没有办法使用flatMap 来实现我想要实现的目标?

【问题讨论】:

  • 您确定要提取元素吗?您的问题似乎是在询问如何将列表列表扁平化为给定 ID 的一个值

标签: scala apache-spark


【解决方案1】:

您可以使用++ 运算符简单地连接两个列表:

val res: RDD[List[String]] = rdd1.join(rdd2)
  .map { case (_, (list1, list2)) => list1 ++ list2 }

可能避免携带可能非常大的List[String] 的更好方法是将RDD 分解为更小的(键值)对,将它们连接起来,然后执行groupByKey

val flatten1: RDD[(String, String)] = rdd1.flatMapValues(identity)
val flatten2: RDD[(String, String)] = rdd2.flatMapValues(identity)
val res: RDD[Iterable[String]] = (flatten1 ++ flatten2).groupByKey.values

【讨论】:

  • 感谢您使用++ 运算符提供的解决方案。我添加了mkString 以获取 RDD[String] (res.map { case (_, (list1, list2)) => (list1 ++ list2).mkString(",") })。我认为您在帖子的第二部分是指(flatten1 ++ flatten2) 而不是(rdd1 ++ rdd2)
  • 我已经尝试过您的建议以避免携带List[String],但我不知道这将如何实现相同的结果,因为我需要在id 上加入我的两个RDD价值。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-02-09
  • 2023-01-28
  • 1970-01-01
  • 2021-08-23
  • 2015-03-24
相关资源
最近更新 更多