【发布时间】:2017-06-10 10:22:16
【问题描述】:
常规的 scala 集合有一个漂亮的 collect 方法,它让我可以使用部分函数一次性执行 filter-map 操作。 sparkDatasets 上是否有等效操作?
我喜欢它有两个原因:
- 语法简洁
- 它将
filter-map样式操作减少到单次传递(尽管在 spark 中我猜有一些优化可以为您发现这些东西)
这里有一个例子来说明我的意思。假设我有一个选项序列,我想提取定义的整数并将其加倍(那些在Some 中的整数):
val input = Seq(Some(3), None, Some(-1), None, Some(4), Some(5))
方法 1 - collect
input.collect {
case Some(value) => value * 2
}
// List(6, -2, 8, 10)
collect 在语法上非常简洁,并且只通过一次。
方法 2 - filter-map
input.filter(_.isDefined).map(_.get * 2)
我可以将这种模式带到 spark 中,因为数据集和数据帧具有类似的方法。
但我不太喜欢这个,因为isDefined 和get 对我来说似乎是代码的味道。有一个隐含的假设,即 map 只接收Somes。编译器无法验证这一点。在更大的示例中,开发人员更难发现这种假设,并且开发人员可能会交换过滤器和映射,例如,不会出现语法错误。
方法 3 - fold* 操作
input.foldRight[List[Int]](Nil) {
case (nextOpt, acc) => nextOpt match {
case Some(next) => next*2 :: acc
case None => acc
}
}
我没有使用足够的 spark 来知道 fold 是否有等价物,所以这可能有点切线。
无论如何,模式匹配、折叠样板和列表的重建都混在一起,很难阅读。
所以总的来说,我发现 collect 语法最好,我希望 spark 有这样的东西。
【问题讨论】:
-
在
RDDs 和Datasets 上定义的collect方法用于实现驱动程序中的数据。尽管没有类似于 Collections APIcollect方法的东西,但您的直觉是正确的:由于这两个操作都是延迟评估的,因此引擎有机会优化操作并将它们链接起来,以便以最大的局部性执行它们。
标签: scala apache-spark apache-spark-dataset