【问题标题】:How to convert spark SchemaRDD into RDD of my case class?如何将 spark SchemaRDD 转换为我的案例类的 RDD?
【发布时间】:2014-11-28 15:41:48
【问题描述】:

在 spark 文档中,很清楚如何从您自己的案例类的 RDD 创建镶木地板文件; (来自文档)

val people: RDD[Person] = ??? // An RDD of case class objects, from the previous example.

// The RDD is implicitly converted to a SchemaRDD by createSchemaRDD, allowing it to be stored using Parquet.
people.saveAsParquetFile("people.parquet")

但不清楚如何转换回来,我们真的想要一个方法readParquetFile 可以做到:

val people: RDD[Person] = sc.readParquestFile[Person](path)

案例类的那些值被定义的地方是那些被方法读取的值。

【问题讨论】:

  • 自从第一次被问到这个问题后有什么更新的解决方案吗?

标签: sql apache-spark parquet


【解决方案1】:

一种简单的方法是提供您自己的转换器(Row) => CaseClass。这有点手动,但如果你知道你在读什么,它应该很简单。

这是一个例子:

import org.apache.spark.sql.SchemaRDD

case class User(data: String, name: String, id: Long)

def sparkSqlToUser(r: Row): Option[User] = {
    r match {
      case Row(time: String, name: String, id: Long) => Some(User(time,name, id))
      case _ => None
    }
}

val parquetData: SchemaRDD = sqlContext.parquetFile("hdfs://localhost/user/data.parquet")

val caseClassRdd: org.apache.spark.rdd.RDD[User] = parquetData.flatMap(sparkSqlToUser)

【讨论】:

    【解决方案2】:

    我想出的需要最少复制和粘贴新类的最佳解决方案如下(不过我仍然希望看到另一个解决方案)

    首先你必须定义你的案例类,和一个(部分)可重用的工厂方法

    import org.apache.spark.sql.catalyst.expressions
    
    case class MyClass(fooBar: Long, fred: Long)
    
    // Here you want to auto gen these functions using macros or something
    object Factories extends java.io.Serializable {
      def longLong[T](fac: (Long, Long) => T)(row: expressions.Row): T = 
        fac(row(0).asInstanceOf[Long], row(1).asInstanceOf[Long])
    }
    

    一些已经可用的样板

    import scala.reflect.runtime.universe._
    val sqlContext = new org.apache.spark.sql.SQLContext(sc)
    import sqlContext.createSchemaRDD
    

    魔法

    import scala.reflect.ClassTag
    import org.apache.spark.sql.SchemaRDD
    
    def camelToUnderscores(name: String) = 
      "[A-Z]".r.replaceAllIn(name, "_" + _.group(0).toLowerCase())
    
    def getCaseMethods[T: TypeTag]: List[String] = typeOf[T].members.sorted.collect {
      case m: MethodSymbol if m.isCaseAccessor => m
    }.toList.map(_.toString)
    
    def caseClassToSQLCols[T: TypeTag]: List[String] = 
      getCaseMethods[T].map(_.split(" ")(1)).map(camelToUnderscores)
    
    def schemaRDDToRDD[T: TypeTag: ClassTag](schemaRDD: SchemaRDD, fac: expressions.Row => T) = {
      val tmpName = "tmpTableName" // Maybe should use a random string
      schemaRDD.registerAsTable(tmpName)
      sqlContext.sql("SELECT " + caseClassToSQLCols[T].mkString(", ") + " FROM " + tmpName)
      .map(fac)
    }
    

    示例使用

    val parquetFile = sqlContext.parquetFile(path)
    
    val normalRDD: RDD[MyClass] = 
      schemaRDDToRDD[MyClass](parquetFile, Factories.longLong[MyClass](MyClass.apply))
    

    另见:

    http://apache-spark-user-list.1001560.n3.nabble.com/Spark-SQL-Convert-SchemaRDD-back-to-RDD-td9071.html

    虽然我没有通过 JIRA 链接找到任何示例或文档。

    【讨论】:

      【解决方案3】:

      在 Spark 1.2.1 中有一个使用 pyspark 将 schema rdd 转换为 rdd 的简单方法。

      sc = SparkContext()  ## create SparkContext
      srdd = sqlContext.sql(sql)
      c = srdd.collect()  ## convert rdd to list
      rdd = sc.parallelize(c)
      

      使用scala必须有类似的方法。

      【讨论】:

      • 这(收集)适用于小型数据收集,但如果您有很多记录,则不适用。
      • 我认为您可能错过了问题的重点,我们想要类似类型提供程序的功能。由于 python 有一个相当弱的类型系统并且是动态的,我怀疑在 python 世界中人们是否真的关心。我们正在构建需要高度稳定的应用程序,因此我们使用具有适当类型系统的语言并需要类型提供程序功能。
      • 对不起,我没抓住重点。我对你的解决方案有点疑惑。sqlContext.sql("SELECT " + caseClassToSQLCols[T].mkString(", ") + " FROM " + tmpName) 返回一个 srdd 对象,所以它没有 map 方法。而 map 方法只是用于将函数应用于每个元素。我找不到您的代码可以将 srdd 转换为 rdd。你能告诉我我是否错过了你的代码中的一些重要内容。谢谢!
      • 另外,我的spark版本是1.2.1。
      【解决方案4】:

      非常粗鲁的尝试。非常不相信这会有不错的表现。当然必须有一个基于宏的替代方案......

      import scala.reflect.runtime.universe.typeOf
      import scala.reflect.runtime.universe.MethodSymbol
      import scala.reflect.runtime.universe.NullaryMethodType
      import scala.reflect.runtime.universe.TypeRef
      import scala.reflect.runtime.universe.Type
      import scala.reflect.runtime.universe.NoType
      import scala.reflect.runtime.universe.termNames
      import scala.reflect.runtime.universe.runtimeMirror
      
      schemaRdd.map(row => RowToCaseClass.rowToCaseClass(row.toSeq, typeOf[X], 0))
      
      object RowToCaseClass {
        // http://dcsobral.blogspot.com/2012/08/json-serialization-with-reflection-in.html
        def rowToCaseClass(record: Seq[_], t: Type, depth: Int): Any = {
          val fields = t.decls.sorted.collect {
            case m: MethodSymbol if m.isCaseAccessor => m
          }
          val values = fields.zipWithIndex.map {
            case (field, i) =>
              field.typeSignature match {
                case NullaryMethodType(sig) if sig =:= typeOf[String] => record(i).asInstanceOf[String]
                case NullaryMethodType(sig) if sig =:= typeOf[Int] => record(i).asInstanceOf[Int]
                case NullaryMethodType(sig) =>
                  if (sig.baseType(typeOf[Seq[_]].typeSymbol) != NoType) {
                    sig match {
                      case TypeRef(_, _, args) =>
                        record(i).asInstanceOf[Seq[Seq[_]]].map {
                          r => rowToCaseClass(r, args(0), depth + 1)
                        }.toSeq
                    }
                  } else {
                    sig match {
                      case TypeRef(_, u, _) =>
                        rowToCaseClass(record(i).asInstanceOf[Seq[_]], sig, depth + 1)
                    }
                  }
              }
          }.asInstanceOf[Seq[Object]]
          val mirror = runtimeMirror(t.getClass.getClassLoader)
          val ctor = t.member(termNames.CONSTRUCTOR).asMethod
          val klass = t.typeSymbol.asClass
          val method = mirror.reflectClass(klass).reflectConstructor(ctor)
          method.apply(values: _*)
        }
      }
      

      【讨论】:

        猜你喜欢
        • 2021-09-27
        • 2023-02-24
        • 2016-08-28
        • 2015-02-27
        • 1970-01-01
        • 1970-01-01
        • 2015-11-11
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多