巧妙处理 Spark RDD 中的 Option[T]

ric*_*din 2 scala optional apache-spark

我正在使用 Apache Spark 的 Scala API 开发一些代码,并且我正在尝试巧妙地解决RDDs之间包含一些Option[T].

假设我们有以下列表

val rdd: RDD[(A, Option[B])] = // Initialization stuff
Run Code Online (Sandbox Code Playgroud)

我们想应用一个转换rdd来获得以下内容

val transformed: RDD[(B, A)]
Run Code Online (Sandbox Code Playgroud)

对于所有Option[B]评估为 的 s Some[B]。我发现这样做的最好方法是应用以下转换链:

val transformed = 
  rdd.filter(_.isDefined)
     .map { case (a, Some(b)) => (b, a) }
Run Code Online (Sandbox Code Playgroud)

我知道如果我使用的是简单的 Scala,List我可以使用以下collect方法:

val transformed = list.collect {
  case (a, Some(b)) => (b, a)
}
Run Code Online (Sandbox Code Playgroud)

如我的这个SO question 所述。

RDD改用Spark s,我有哪种选择?

zer*_*323 5

RDD提供等效于的转换collect Iterable.collect

import org.apache.spark.rdd.RDD

val rdd = sc.parallelize(Seq((1L, None), (2L, Some("a"))))
val transformed: RDD[(String, Long)] = rdd.collect {
  case (a, Some(b)) => (b, a)
}

transformed.count
// Long = 1 

transformed.first
// (String, Long) = (a,2)
Run Code Online (Sandbox Code Playgroud)


Jea*_*art 5

您可以使用flatMap

rdd.flatMap {
   case (a, Some(b)) => Some(b, a)
   case _ => None
}
Run Code Online (Sandbox Code Playgroud)