【问题标题】:How do I change the schema on a Spark Dataset如何更改 Spark 数据集上的架构
【发布时间】:2017-12-23 02:00:18
【问题描述】:

当我在 Spark 2 中检索数据集时,使用 select 语句,底层列会继承查询列的数据类型。

val ds1 = spark.sql("select 1 as a, 2 as b, 'abd' as c")

ds1.printSchema()
root
 |-- a: integer (nullable = false)
 |-- b: integer (nullable = false)
 |-- c: string (nullable = false)

现在,如果我将其转换为案例类,它会正确转换值,但底层架构仍然是错误的。

case class abc(a: Double, b: Double, c: String)
val ds2 = ds1.as[abc]
ds2.printSchema()
root
 |-- a: integer (nullable = false)
 |-- b: integer (nullable = false)
 |-- c: string (nullable = false)

ds2.collect
res18: Array[abc] = Array(abc(1.0,2.0,abd))

当我创建第二个数据集时,我“应该”能够指定要使用的编码器,但 scala 似乎忽略了这个参数(这是一个 BUG 吗?):

val abc_enc = org.apache.spark.sql.Encoders.product[abc]

val ds2 = ds1.as[abc](abc_enc)

ds2.printSchema
root
 |-- a: integer (nullable = false)
 |-- b: integer (nullable = false)
 |-- c: string (nullable = false)

所以我能看到的唯一方法是简单地做到这一点,没有非常复杂的映射是使用 createDataset,但这需要对底层对象进行收集,所以它并不理想。

val ds2 = spark.createDataset(ds1.as[abc].collect)

【问题讨论】:

    标签: scala apache-spark encoder apache-spark-dataset


    【解决方案1】:

    你可以简单地在columns上使用cast方法

    import sqlContext.implicits._
    val ds2 = ds1.select($"a".cast(DoubleType), $"a".cast(DoubleType), $"c")
    ds2.printSchema()
    

    你应该有

    root
     |-- a: double (nullable = false)
     |-- a: double (nullable = false)
     |-- c: string (nullable = false)
    

    【讨论】:

    • 答案对您没有帮助吗?如果确实如此,那么考虑接受和支持:)
    【解决方案2】:

    您还可以在使用 sql 查询进行选择时强制转换列,如下所示

    import spark.implicits._
    
    val ds = Seq((1,2,"abc"),(1,2,"abc")).toDF("a", "b","c").createOrReplaceTempView("temp")
    
    val ds1 = spark.sql("select cast(a as Double) , cast (b as Double), c from temp")
    
    ds1.printSchema()
    

    这有架构为

    root
     |-- a: double (nullable = false)
     |-- b: double (nullable = false)
     |-- c: string (nullable = true)
    

    现在您可以转换为带有案例类的数据集

    case class abc(a: Double, b: Double, c: String)
    
    val ds2 = ds1.as[abc]
    ds2.printSchema()
    

    现在有所需的架构

    root
     |-- a: double (nullable = false)
     |-- b: double (nullable = false)
     |-- c: string (nullable = true)
    

    希望这会有所帮助!

    【讨论】:

      【解决方案3】:

      好的,我想我已经以更好的方式解决了这个问题。

      当我们创建一个新数据集时,我们可以直接引用数据集的rdd,而不是使用collect。

      所以不是

      val ds2 = spark.createDataset(ds1.as[abc].collect)
      

      我们使用:

      val ds2 = spark.createDataset(ds1.as[abc].rdd)
      
      ds2.printSchema
      root
       |-- a: double (nullable = false)
       |-- b: double (nullable = false)
       |-- c: string (nullable = true)
      

      这使得惰性求值保持不变,但允许新数据集使用 abc 案例类的 Encoder,后续架构在我们使用它创建新表时会反映这一点。

      【讨论】:

        【解决方案4】:

        这是 Spark API 中的一个未解决问题(请查看此票证 SPARK-17694

        所以你需要做的是做一个额外的显式转换。像这样的东西应该可以工作:

        ds1.as[abc].map(x => x : abc)
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2011-03-25
          • 2017-09-22
          • 2014-11-03
          • 1970-01-01
          • 2020-09-23
          • 2021-12-15
          • 1970-01-01
          相关资源
          最近更新 更多