【问题标题】:Polymorphism with Spark / Scala, Datasets and case classesSpark / Scala、数据集和案例类的多态性
【发布时间】:2017-06-22 19:52:01
【问题描述】:

我们将 Spark 2.x 与 Scala 一起用于具有 13 种不同 ETL 操作的系统。其中 7 个相对简单,每个都由一个域类驱动,主要区别在于这个类和负载处理方式的一些细微差别。

加载类的简化版本如下,为了这个示例的目的,假设要加载 7 个披萨浇头,这里是意大利辣香肠:

object LoadPepperoni {
  def apply(inputFile: Dataset[Row],
            historicalData: Dataset[Pepperoni],
            mergeFun: (Pepperoni, PepperoniRaw) => Pepperoni): Dataset[Pepperoni] = {
    val sparkSession = SparkSession.builder().getOrCreate()
    import sparkSession.implicits._

    val rawData: Dataset[PepperoniRaw] = inputFile.rdd.map{ case row : Row =>
      PepperoniRaw(
          weight = row.getAs[String]("weight"),
          cost = row.getAs[String]("cost")
        )
    }.toDS()

    val validatedData: Dataset[PepperoniRaw] = ??? // validate the data

    val dedupedRawData: Dataset[PepperoniRaw] = ??? // deduplicate the data

    val dedupedData: Dataset[Pepperoni] = dedupedRawData.rdd.map{ case datum : PepperoniRaw =>
        Pepperoni( value = ???, key1 = ???, key2 = ??? )
    }.toDS()

    val joinedData = dedupedData.joinWith(historicalData,
      historicalData.col("key1") === dedupedData.col("key1") && 
        historicalData.col("key2") === dedupedData.col("key2"),
      "right_outer"
    )

    joinedData.map { case (hist, delta) =>
      if( /* some condition */) {
        hist.copy(value = /* some transformation */)
      }
    }.flatMap(list => list).toDS()
  }
}

换句话说,该类对数据执行一系列操作,这些操作大多相同且顺序始终相同,但每个顶部可能略有不同,从“原始”到“域”的映射和合并函数。

要对 7 种配料(即蘑菇、奶酪等)执行此操作,我不想简单地复制/粘贴类并更改所有名称,因为结构和逻辑对所有负载都是通用的。相反,我宁愿定义一个具有泛型类型的泛型“Load”类,如下所示:

object Load {
  def apply[R,D](inputFile: Dataset[Row],
            historicalData: Dataset[D],
            mergeFun: (D, R) => D): Dataset[D] = {
    val sparkSession = SparkSession.builder().getOrCreate()
    import sparkSession.implicits._

    val rawData: Dataset[R] = inputFile.rdd.map{ case row : Row =>
...

并且对于每个特定于类的操作,例如从“原始”到“域”的映射或合并,都有一个实现细节的特征或抽象类。这将是一个典型的依赖注入/多态模式。

但我遇到了一些问题。从 Spark 2.x 开始,只为原生类型和案例类提供编码器,没有办法将一个类一般地标识为案例类。因此推断的 toDS() 和其他隐式功能在使用泛型类型时不可用。

同样在this related question of mine中提到,使用泛型时case类copy方法也不可用。

我研究了 Scala 和 Haskell 常见的其他设计模式,例如类型类或 ad-hoc 多态性,但障碍是 Spark 数据集基本上只适用于无法抽象定义的案例类。

这似乎是 Spark 系统中的常见问题,但我无法找到解决方案。任何帮助表示赞赏。

【问题讨论】:

    标签: scala apache-spark design-patterns polymorphism


    【解决方案1】:

    启用.toDS 的隐式转换为:

    implicit def rddToDatasetHolder[T](rdd: RDD[T])(implicit arg0: Encoder[T]): DatasetHolder[T]
    

    (来自https://spark.apache.org/docs/latest/api/scala/index.html#org.apache.spark.sql.SQLImplicits

    您完全正确,因为您已经使您的 apply 方法通用,因此 Encoder[T] 的范围内没有隐式值,因此这种转换不会发生。但是你可以简单地接受一个作为隐式参数!

    object Load {
      def apply[R,D](inputFile: Dataset[Row],
                historicalData: Dataset[D],
                mergeFun: (D, R) => D)(implicit enc: Encoder[D]): Dataset[D] = {
    ...
    

    然后在您调用负载时,具有特定类型,它应该能够找到该类型的编码器。请注意,您还必须在调用上下文中使用import sparkSession.implicits._

    编辑:类似的方法是通过限制类型 (apply[R, D <: Product]) 并接受隐式 JavaUniverse.TypeTag[D] 作为参数来启用隐式 newProductEncoder[T <: Product](implicit arg0: scala.reflect.api.JavaUniverse.TypeTag[T]): Encoder[T]

    【讨论】:

    • 谢谢你,这让我大开眼界,让我试一试。我在 aggregateByKey() 和 copy() 上遇到了类似的错误(请参阅我帖子中的相关问题链接),是否有任何类似的魔法可以将所需的实现纳入范围?
    • 更一般地说,该方法是跟踪丢失的特定隐式并将其作为参数一直到类型完全指定的调用站点。在 aggregateByKey 的情况下,您似乎需要一个 ClassTag[D] (spark.apache.org/docs/latest/api/scala/…)。我会单独看看你关于copy的问题,如果我有什么好的想法,我会在那里回复
    猜你喜欢
    • 2021-12-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多