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,我有哪种选择?
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)
您可以使用flatMap:
rdd.flatMap {
case (a, Some(b)) => Some(b, a)
case _ => None
}
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
2328 次 |
| 最近记录: |