【问题标题】:Spark MergeSchema on parquet columns镶木地板上的 Spark MergeSchema
【发布时间】:2020-04-19 11:51:08
【问题描述】:

对于模式演变,Mergeschema 可以在 Spark 中用于 Parquet 文件格式,我对此有以下说明

这是否仅支持 Parquet 文件格式或任何其他文件格式,如 csv、txt 文件。

如果在两者之间添加了新的附加列,我了解 Mergeschema 会将这些列移到最后。

如果列顺序受到干扰,那么 Mergeschema 是否会在创建时将列对齐以正确顺序,还是我们需要通过选择所有列手动执行此操作。

评论更新: 例如,如果我有如下架构并创建如下表 - spark.sql("CREATE TABLE emp USING DELTA LOCATION '****'") empid,empname,salary====> 001,ABC,10000 和第二天如果我得到以下格式 empid,empage,empdept,empname,salary====> 001,30,XYZ,ABC,10000。

empid,empname,salary columns之后是否会添加新列-empage, empdept?

【问题讨论】:

    标签: scala azure apache-spark databricks


    【解决方案1】:

    问: 1. 这是否只支持 Parquet 文件格式或任何其他文件格式,如 csv、txt 文件。 2. 如果列顺序受到干扰,Mergeschema 是否会在创建时将列对齐到正确的顺序,还是我们需要通过选择所有列来手动执行此操作


    只有 parquet 支持 AFAIK Merge 架构,其他格式(如 csv、txt)不支持。

    Mergeschema (spark.sql.parquet.mergeSchema) 将按正确的顺序对齐列,即使它们是分布的。

    parquet schema-merging 上的 spark 文档示例:

    import spark.implicits._
    
    // Create a simple DataFrame, store into a partition directory
    val squaresDF = spark.sparkContext.makeRDD(1 to 5).map(i => (i, i * i)).toDF("value", "square")
    squaresDF.write.parquet("data/test_table/key=1")
    
    // Create another DataFrame in a new partition directory,
    // adding a new column and dropping an existing column
    val cubesDF = spark.sparkContext.makeRDD(6 to 10).map(i => (i, i * i * i)).toDF("value", "cube")
    cubesDF.write.parquet("data/test_table/key=2")
    
    // Read the partitioned table
    val mergedDF = spark.read.option("mergeSchema", "true").parquet("data/test_table")
    mergedDF.printSchema()
    
    // The final schema consists of all 3 columns in the Parquet files together
    // with the partitioning column appeared in the partition directory paths
    // root
    //  |-- value: int (nullable = true)
    //  |-- square: int (nullable = true)
    //  |-- cube: int (nullable = true)
    //  |-- key: int (nullable = true)
    

    更新:你在评论框中给出的真实例子......


    Q : 之后是否会添加新的列 - empage, empdept empid,empname,salary columns?


    答案:是的 在 EMPID,EMPNAME,SALARY 之后添加了 EMPAGE,EMPDEPT,然后是您的日期列。

    查看完整示例。

    package examples
    
    import org.apache.log4j.Level
    import org.apache.spark.sql.SaveMode
    
    
    object CSVDataSourceParquetSchemaMerge extends App {
      val logger = org.apache.log4j.Logger.getLogger("org")
      logger.setLevel(Level.WARN)
    
      import org.apache.spark.sql.SparkSession
    
      val spark = SparkSession.builder().appName("CSVParquetSchemaMerge")
        .master("local")
        .getOrCreate()
    
    
      import spark.implicits._
    
      val csvDataday1 = spark.sparkContext.parallelize(
        """
          |empid,empname,salary
          |001,ABC,10000
        """.stripMargin.lines.toList).toDS()
      val csvDataday2 = spark.sparkContext.parallelize(
        """
          |empid,empage,empdept,empname,salary
          |001,30,XYZ,ABC,10000
        """.stripMargin.lines.toList).toDS()
    
      val frame = spark.read.option("header", true).option("inferSchema", true).csv(csvDataday1)
    
      println("first day data ")
      frame.show
      frame.write.mode(SaveMode.Overwrite).parquet("data/test_table/day=1")
      frame.printSchema
    
      val frame1 = spark.read.option("header", true).option("inferSchema", true).csv(csvDataday2)
      frame1.write.mode(SaveMode.Overwrite).parquet("data/test_table/day=2")
      println("Second day data ")
    
      frame1.show(false)
      frame1.printSchema
    
      // Read the partitioned table
      val mergedDF = spark.read.option("mergeSchema", "true").parquet("data/test_table")
      println("Merged Schema")
      mergedDF.printSchema
      println("Merged Datarame where EMPAGE,EMPDEPT WERE ADDED AFER EMPID,EMPNAME,SALARY followed by your day column")
      mergedDF.show(false)
    
    
    }
    
    
    

    结果:

    first day data 
    +-----+-------+------+
    |empid|empname|salary|
    +-----+-------+------+
    |    1|    ABC| 10000|
    +-----+-------+------+
    
    root
     |-- empid: integer (nullable = true)
     |-- empname: string (nullable = true)
     |-- salary: integer (nullable = true)
    
    Second day data 
    +-----+------+-------+-------+------+
    |empid|empage|empdept|empname|salary|
    +-----+------+-------+-------+------+
    |1    |30    |XYZ    |ABC    |10000 |
    +-----+------+-------+-------+------+
    
    root
     |-- empid: integer (nullable = true)
     |-- empage: integer (nullable = true)
     |-- empdept: string (nullable = true)
     |-- empname: string (nullable = true)
     |-- salary: integer (nullable = true)
    
    Merged Schema
    root
     |-- empid: integer (nullable = true)
     |-- empname: string (nullable = true)
     |-- salary: integer (nullable = true)
     |-- empage: integer (nullable = true)
     |-- empdept: string (nullable = true)
     |-- day: integer (nullable = true)
    
    Merged Datarame where EMPAGE,EMPDEPT WERE ADDED AFER EMPID,EMPNAME,SALARY followed by your day column
    +-----+-------+------+------+-------+---+
    |empid|empname|salary|empage|empdept|day|
    +-----+-------+------+------+-------+---+
    |1    |ABC    |10000 |30    |XYZ    |2  |
    |1    |ABC    |10000 |null  |null   |1  |
    +-----+-------+------+------+-------+---+
    
    

    目录树:

    【讨论】:

    • 例如,如果我有如下架构并创建如下表 - spark.sql("CREATE TABLE emp USING DELTA LOCATION '****'") empid,empname,salary=== => 001,ABC,10000 和第二天如果我得到以下格式 empid,empage,empdept,empname,salary====> 001,30,XYZ,ABC,10000 是否会在 empid 之后添加新列 - empage, empdept ,empname,salary 列?
    • 为什么不以上述方式与您的示例一起尝试呢?
    • 是的,尝试了一个示例,它正在将列移到最后..只是想检查您是否已经知道..
    猜你喜欢
    • 2017-03-17
    • 2016-07-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-20
    • 2019-06-02
    • 2016-07-07
    • 2016-06-16
    相关资源
    最近更新 更多