使用map函数检查RDD元素是否在另一个元素中

Jos*_*yia 0 closures scala apache-spark

我是Spark的新手,对封口感到疑惑.
我有两个RDD,一个包含ID和值列表,另一个包含所选ID列表.
使用map,我想增加元素的值,如果另一个RDD包含它的ID,就像这样.

val ids = sc.parallelize(List(1,2,10,5))
val vals = sc.parallelize(List((1, 0), (2, 0), (3,0), (4,0)))
vals.map( v => {
  if(ids.collect().contains(v._1)){
    (v._1, 1)
  } 
 })
Run Code Online (Sandbox Code Playgroud)

然而,工作挂起并永远不会完成.这样做的正确方法是什么,谢谢你的帮助!

Tza*_*har 5

您的实现尝试在ids用于映射另一个的闭包内使用一个RDD() - 这在Spark应用程序中是不允许的:在闭包中使用的任何东西都必须是可序列化的(并且最好是小的),因为它将被序列化并发送到每个工人.

leftOuterJoin,这些RDDS之间应该得到你想要的东西:

val ids = sc.parallelize(List(1,2,10,5))
val vals = sc.parallelize(List((1, 0), (2, 0), (3,0), (4,0)))
val result = vals
        .leftOuterJoin(ids.keyBy(i => i))
        .mapValues({ 
            case (v, Some(matchingId)) => v + 1  // increase value if match found
            case (v, None) => v                  // leave value as-is otherwise
        }) 
Run Code Online (Sandbox Code Playgroud)

leftOuterJoin需要两个键值RDDS,因此我们人为地提取从一个密钥ids使用标识功能RDD.然后我们将每个结果记录的映射(id: Int, (value: Int, matchingId: Option[Int]))到v或v + 1.

通常,您应该始终致力于最大限度地减少collect使用Spark时的操作,因为此类操作会将数据从分布式群集移回驱动程序应用程序.