【发布时间】:2015-12-22 23:32:59
【问题描述】:
因此,在面向对象的世界中度过了多年,始终考虑到代码重用、设计模式和最佳实践,我发现自己在 Spark 世界中的代码组织和代码重用方面有些挣扎。
如果我尝试以可重用的方式编写代码,它几乎总是会带来性能成本,我最终会将其重写为最适合我的特定用例的任何内容。这个常量“为这个特定用例编写最佳的东西”也影响代码组织,因为当“它真的属于一起”时,将代码分成不同的对象或模块是很困难的,因此我最终得到的“上帝”对象很少包含 long复杂的转换链。事实上,我经常认为,如果我回顾一下我现在在面向对象领域工作时正在编写的大部分 Spark 代码,我会畏缩不前,并认为它是“意大利面条式代码”。
我在网上冲浪,试图找到某种与面向对象世界的最佳实践等效的方法,但运气不佳。我可以找到一些函数式编程的“最佳实践”,但 Spark 只是添加了一个额外的层,因为性能是这里的一个主要因素。
所以我的问题是,你们中的任何一位 Spark 专家是否找到了一些可以推荐的编写 Spark 代码的最佳实践?
编辑
正如评论中所写,我实际上并不希望有人发布关于如何解决这个问题的答案,而是我希望这个社区中的某个人遇到过一些 Martin Fowler 类型,谁曾在某处写过一些关于如何解决 Spark 世界中代码组织问题的文章或博客文章。
@DanielDarabos 建议我可以举一个代码组织和性能冲突的例子。虽然我发现我在日常工作中经常遇到这个问题,但我发现将其归结为一个好的最小示例有点困难;)但我会尝试。
在面向对象的世界中,我是单一职责原则的忠实拥护者,所以我会确保我的方法只对一件事负责。它使它们可重用且易于测试。因此,如果我必须计算列表中某些数字的总和(匹配某些标准)并且我必须计算相同数字的平均值,我肯定会创建两种方法 - 一种计算总和,另一种计算总和计算了平均值。像这样:
def main(implicit args: Array[String]): Unit = {
val list = List(("DK", 1.2), ("DK", 1.4), ("SE", 1.5))
println("Summed weights for DK = " + summedWeights(list, "DK")
println("Averaged weights for DK = " + averagedWeights(list, "DK")
}
def summedWeights(list: List, country: String): Double = {
list.filter(_._1 == country).map(_._2).sum
}
def averagedWeights(list: List, country: String): Double = {
val filteredByCountry = list.filter(_._1 == country)
filteredByCountry.map(_._2).sum/ filteredByCountry.length
}
我当然可以继续尊重 Spark 中的 SRP:
def main(implicit args: Array[String]): Unit = {
val df = List(("DK", 1.2), ("DK", 1.4), ("SE", 1.5)).toDF("country", "weight")
println("Summed weights for DK = " + summedWeights(df, "DK")
println("Averaged weights for DK = " + averagedWeights(df, "DK")
}
def avgWeights(df: DataFrame, country: String, sqlContext: SQLContext): Double = {
import org.apache.spark.sql.functions._
import sqlContext.implicits._
val countrySpecific = df.filter('country === country)
val summedWeight = countrySpecific.agg(avg('weight))
summedWeight.first().getDouble(0)
}
def summedWeights(df: DataFrame, country: String, sqlContext: SQLContext): Double = {
import org.apache.spark.sql.functions._
import sqlContext.implicits._
val countrySpecific = df.filter('country === country)
val summedWeight = countrySpecific.agg(sum('weight))
summedWeight.first().getDouble(0)
}
但是因为我的df 可能包含数十亿行,我宁愿不必执行filter 两次。事实上,性能与 EMR 成本直接相关,所以我真的不希望这样。为了克服它,我因此决定违反 SRP 并简单地将两个函数合二为一,并确保我在国家过滤的DataFrame 上调用persist,如下所示:
def summedAndAveragedWeights(df: DataFrame, country: String, sqlContext: SQLContext): (Double, Double) = {
import org.apache.spark.sql.functions._
import sqlContext.implicits._
val countrySpecific = df.filter('country === country).persist(StorageLevel.MEMORY_AND_DISK_SER)
val summedWeights = countrySpecific.agg(sum('weight)).first().getDouble(0)
val averagedWeights = summedWeights / countrySpecific.count()
(summedWeights, averagedWeights)
}
现在,这个例子当然是对现实生活中遇到的事情的极大简化。在这里,我可以简单地通过过滤和持久化df 在将它交给 sum 和 avg 函数(这也将是更多的 SRP)来解决它,但在现实生活中可能会有许多中间计算需要一次又一次地进行下去。换句话说,这里的filter 函数只是试图制作一个简单 的例子,说明一些可以从持久化中受益的东西。事实上,我认为调用persist 是这里的关键字。调用persist 将大大加快我的工作,但代价是我必须紧密耦合依赖于持久化DataFrame 的所有代码——即使它们在逻辑上是分开的。
【问题讨论】:
-
有什么特别的语言吗?我根本不是专家,但对于 Java 和 Scala,我认为没有理由不按照自己的标准构建代码。 Databriks 参考应用程序 (github.com/databricks/reference-apps/tree/master/timeseries) 是构建 Spark 项目的一个非常好的开始。希望对您有所帮助!
-
我探索了不同的方法,并在我第一次了解数据集时拥有您所说的意大利面条代码。然后我考虑如何对我正在使用的数据进行分类。它是如何变异的,等等。从那里开始,经典的软件设计模式往往对我有用。
-
另外 - 我没有看到业界对如何在可重用的分布式环境中编写高效、可扩展的代码达成共识。这些模式通常与数据高度耦合,因此您必须努力使用商定的标准创建数据。对于某些问题,这永远不够高效。
-
我完全理解这个问题,但我认为它涵盖了除重复之外的所有 Stack Overflow 关闭原因。你觉得另一种表达方式如何?也许展示一个代码组织和性能是冲突目标的案例的最小示例,并询问如何解决该冲突。我认为它可以很好地作为当前问题内容的补充,并可以给出具体的答案。
-
@DanielDarabos 我理解你的评论,我正在想一个很好的例子来回答这个问题。
标签: apache-spark functional-programming code-organization