【问题标题】:Perform a nested for loop with RDD.map() in Scala在 Scala 中使用 RDD.map() 执行嵌套的 for 循环
【发布时间】:2017-10-12 11:36:27
【问题描述】:

我对 Spark 和 Scala 比较陌生,并且有 Java 背景。我在 haskell 中做过一些编程,所以对函数式编程并不完全陌生。

我正在尝试完成某种形式的嵌套 for 循环。我有一个 RDD,我想根据 RDD 中的每两个元素对其进行操作。伪代码(类似 java)如下所示:

// some RDD named rdd is available before this
List list = new ArrayList();
for(int i = 0; i < rdd.length; i++){
   list.add(rdd.get(i)._1);
   for(int j = 0; j < rdd.length; j++){
      if(rdd.get(i)._1 == rdd.get(j)._1){
         list.add(rdd.get(j)._1);
      }
   }
}
// Then now let ._1 of the rdd be this list

我的 scala 解决方案(不起作用)如下所示:

  val aggregatedTransactions = joinedTransactions.map( f => {
     var list = List[Any](f._2._1)
     val filtered = joinedTransactions.filter(t => f._1 == t._1)

     for(i <- filtered){
       list ::= i._2._1
     }

     (f._1, list, f._2._2)
  })

如果两个项目的 ._1 相等,我正在尝试将项目 _2._1 放入列表中。 我知道我不能在另一个地图功能中执行任何过滤器或地图功能。我已经读到您可以通过连接实现类似的目标,但我不知道如何将这些项目实际放入列表或任何可用作列表的结构中。

如何使用 RDD 实现这样的效果?

【问题讨论】:

  • 我认为您需要更准确地说明您想要实现的目标(即我认为 Java 代码与您声明的意图不符)。首先,您为什么不使用案例类来定义您正在使用的对象?
  • 如果您是第一次使用 scala,我强烈建议您花一些时间玩 scala,尤其是 scala 集合。希望这有帮助
  • 我不能为此使用 scala 集合,因为集合无法序列化,因此会在 spark 系统上引发错误(由于垃圾收集器超时运行..)。这确实是我的第一次尝试。

标签: scala apache-spark rdd


【解决方案1】:

假设您的输入对于某些类型A, B 具有RDD[(A, (A, B))] 的形式,并且预期结果应该具有RDD[A] 的形式——而不是列表(因为我们希望保持数据分布)——这似乎可以你需要什么:

rdd.join(rdd.values).keys

详情:

很难理解确切的输入和预期输出,因为两者的数据结构(类型)都没有明确说明,代码示例也没有很好地解释需求。所以我会做一些假设,希望对你的具体情况有所帮助。

对于完整的示例,我假设:

  • 输入 RDD 的类型为 RDD[(Int, (Int, Int))]
  • 预期的输出格式为RDD[Int],并且包含大量重复项 - 如果原始 RDD 多次具有“键”X,则每次出现 X 时,每个匹配项(在 ._2._1 中)将出现一次键

如果我们正在尝试解决这种情况 - 这个join 会解决它:

// Some sample data, assuming all ints
val rdd = sc.parallelize(Seq(
  (1, (1, 5)),
  (1, (2, 5)),
  (2, (1, 5)),
  (3, (4, 5))
))

// joining the original RDD with an RDD of the "values" -
// so the joined RDD will have "._2._1" as key
// then we get the keys only, because they equal the values anyway
val result: RDD[Int] = rdd.join(rdd.values).keys

// result is a key-value RDD with the original keys as keys, and a list of matching _2._1
println(result.collect.toList) // List(1, 1, 1, 1, 2)

【讨论】:

  • 我现在意识到我应该更具体地陈述我的问题。我想要实现的是一个包含聚合边缘的 RDD。 ._1 是我想加入他们的关键。结构是(id,(src,dest))。因此,我想获得一个来源列表,同时每个 RDD“行”只保留一个 id。所以: (id, (list, dest)) 在查看您的代码时,我仍然不完全确定如何实现这一点。你能解释一下吗?
  • 这听起来与这篇文章所描述的非常不同(或者如果 Java 代码有效的话会做什么......)所以我建议你发布一个带有信息的 new 问题在此评论中,加上示例输入和预期输出。然后有人(也许是我!)可以提供帮助。
猜你喜欢
  • 2022-06-15
  • 1970-01-01
  • 1970-01-01
  • 2012-09-09
  • 1970-01-01
  • 2021-12-11
  • 2019-05-11
  • 2014-01-12
  • 2015-01-28
相关资源
最近更新 更多