【问题标题】:How to match Dataframe column names to Scala case class attributes?如何将 Dataframe 列名与 Scala 案例类属性匹配?
【发布时间】:2015-12-08 17:05:00
【问题描述】:

此示例中来自 spark-sql 的列名来自 case class Person

case class Person(name: String, age: Int)

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")

https://spark.apache.org/docs/1.1.0/sql-programming-guide.html

但是,在许多情况下,参数名称可能会更改。如果文件尚未更新以反映更改,这将导致找不到列。

如何指定合适的映射?

我在想这样的事情:

  val schema = StructType(Seq(
    StructField("name", StringType, nullable = false),
    StructField("age", IntegerType, nullable = false)
  ))


  val ps: Seq[Person] = ???

  val personRDD = sc.parallelize(ps)

  // Apply the schema to the RDD.
  val personDF: DataFrame = sqlContext.createDataFrame(personRDD, schema)

【问题讨论】:

  • 不幸的是,不清楚你想要什么。 1.用任意名字写拼花? 2. 之后更改拼花列名称? 3. 读取任意列名的 parquet 并将其“匹配”/映射到 case 类的字段?
  • @MartinSenne 怎么样?我想手动设置列名并将案例类参数映射到这些列。
  • 但您打算让它们自动匹配?
  • @MartinSenne 请对此进行扩展。就像我说的我想手动匹配。

标签: scala apache-spark apache-spark-sql parquet


【解决方案1】:

基本上,您需要做的所有映射都可以通过DataFrame.select(...) 来实现。 (在这里,我假设不需要进行类型转换。) 鉴于前向和后向映射作为映射,基本部分是

val mapping = from.map{ (x:(String, String)) => personsDF(x._1).as(x._2) }.toArray
// personsDF your original dataframe  
val mappedDF = personsDF.select( mapping: _* )

其中 mapping 是一个带有别名的 Columns 数组。

示例代码

object Example {   

  import org.apache.spark.rdd.RDD
  import org.apache.spark.{SparkContext, SparkConf}

  case class Person(name: String, age: Int)

  object Mapping {
    val from = Map("name" -> "a", "age" -> "b")
    val to = Map("a" -> "name", "b" -> "age")
  }

  def main(args: Array[String]) : Unit = {
    // init
    val conf = new SparkConf()
      .setAppName( "Example." )
      .setMaster( "local[*]")

    val sc = new SparkContext(conf)
    val sqlContext = new SQLContext(sc)
    import sqlContext.implicits._

    // create persons
    val persons = Seq(Person("bob", 35), Person("alice", 27))
    val personsRDD = sc.parallelize(persons, 4)
    val personsDF = personsRDD.toDF

    writeParquet( personsDF, "persons.parquet", sc, sqlContext)

    val otherPersonDF = readParquet( "persons.parquet", sc, sqlContext )
  }

  def writeParquet(personsDF: DataFrame, path:String, sc: SparkContext, sqlContext: SQLContext) : Unit = {
    import Mapping.from

    val mapping = from.map{ (x:(String, String)) => personsDF(x._1).as(x._2) }.toArray

    val mappedDF = personsDF.select( mapping: _* )
    mappedDF.write.parquet("/output/path.parquet") // parquet with columns "a" and "b"
  }

  def readParquet(path: String, sc: SparkContext, sqlContext: SQLContext) : Unit = {
    import Mapping.to
    val df = sqlContext.read.parquet(path) // this df has columns a and b

    val mapping = to.map{ (x:(String, String)) => df(x._1).as(x._2) }.toArray
    df.select( mapping: _* )
  }
}

备注

如果您需要将数据帧转换回 RDD[Person],那么

val rdd : RDD[Row] = personsDF.rdd
val personsRDD : Rdd[Person] = rdd.map { r: Row => 
  Person( r.getAs("person"), r.getAs("age") )
}

替代方案

也看看How to convert spark SchemaRDD into RDD of my case class?

【讨论】:

  • 不错的方法。您认为这会影响性能吗,还是应该不是一个因素,因为它在内部管道中编译和优化过一次?
  • 我假设是后者。首先,因为有 Catalyst 优化/编译。其次,选择(带有别名)似乎不是昂贵的操作。不过,有兴趣看到一些性能测量......
  • @BAR 我们可以在这里使用 jave 吗?对于给出的例子?在java数据集中选择方法没有能力获取地图?
猜你喜欢
  • 1970-01-01
  • 2014-08-12
  • 1970-01-01
  • 1970-01-01
  • 2017-02-25
  • 2023-03-11
  • 2020-10-04
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多