Spark 和 Scala 中数据框的转换模式

Mas*_*cci 5 scala apache-spark apache-spark-sql spark-dataframe

我想转换数据框的模式以使用 Spark 和 Scala 更改某些列的类型。

具体来说,我试图使用 as[U] 函数,其描述为:“返回一个新的数据集,其中每个记录都已映射到指定的类型。用于映射列的方法取决于 U 的类型”

原则上这正是我想要的,但我无法让它工作。

这是一个来自https://github.com/apache/spark/blob/master/sql/core/src/test/scala/org/apache/spark/sql/DatasetSuite.scala的简单示例



    // definition of data
    val data = Seq(("a", 1), ("b", 2)).toDF("a", "b")

Run Code Online (Sandbox Code Playgroud)

正如预期的那样,数据模式是:

    根
     |-- a: 字符串 (nullable = true)
     |-- b:整数(可为空 = false)
    

我想将列“b”转换为 Double。所以我尝试以下操作:



    import session.implicits._;

    println(" --------------------------- Casting using (String Double)")

    val data_TupleCast=data.as[(String, Double)]
    data_TupleCast.show()
    data_TupleCast.printSchema()

    println(" --------------------------- Casting using ClassData_Double")

    case class ClassData_Double(a: String, b: Double)

    val data_ClassCast= data.as[ClassData_Double]
    data_ClassCast.show()
    data_ClassCast.printSchema()

Run Code Online (Sandbox Code Playgroud)

据我了解 as[u] 的定义,新的 DataFrame 应具有以下架构

    根
     |-- a: 字符串 (nullable = true)
     |-- b: double (nullable = false)

但输出是

     --------------------------- 使用 (String Double) 进行转换
    +---+---+
    | 一个| 乙|
    +---+---+
    | 一个| 1|
    | 乙| 2|
    +---+---+

    根
     |-- a: 字符串 (nullable = true)
     |-- b:整数(可为空 = false)

     --------------------------- 使用 ClassData_Double 进行投射
    +---+---+
    | 一个| 乙|
    +---+---+
    | 一个| 1|
    | 乙| 2|
    +---+---+

    根
     |-- a: 字符串 (nullable = true)
     |-- b:整数(可为空 = false)

这表明列“b”没有被强制转换为双倍。

关于我做错了什么的任何提示?

顺便说一句:我知道上一篇文章“如何更改 Spark SQL 的 DataFrame 中的列类型?” (请参阅如何更改 Spark SQL 的 DataFrame 中的列类型?)。我知道我可以一次更改一个列的类型,但我正在寻找一种更通用的解决方案,可以一次性更改整个数据的架构(并且我试图在此过程中了解 Spark)。

Gle*_*olt 5

好吧,由于函数是链接的并且 Spark 执行惰性求值,因此它实际上确实会一次性更改整个数据的架构,即使您确实将其写为更改一列,如下所示:

import spark.implicits._

df.withColumn("x", 'x.cast(DoubleType)).withColumn("y", 'y.cast(StringType))...
Run Code Online (Sandbox Code Playgroud)

作为替代方案,我认为您可以map一次性完成您的演员表,例如:

df.map{t => (t._1, t._2.asInstanceOf[Double], t._3.asInstanceOf[], ...)}
Run Code Online (Sandbox Code Playgroud)