【问题标题】:How to detect change in CSV file schema in Spark如何在 Spark 中检测 CSV 文件模式的变化
【发布时间】:2020-04-25 17:29:32
【问题描述】:

如果传入的 CSV 文件中的架构发生变化,我们如何在 spark 中处理?

假设在第 1 天, 我得到了一个 csv 文件,其架构和数据如下,

FirstName LastName Age
Sagar     Patro    26
Akash     Nayak    22
Amar      Kumar    18

在第 10 天, 我传入的 CSV 文件架构发生了变化,如下所示

FirstName LastName Mobile     Age 
Sagar     Patro    8984159475 26  
Akash     Nayak    9040988503 22  
Amar      Kumar    9337856871 18  

我的第 1 项要求,

我想知道,我传入的 CSV 文件的架构是否有任何变化。

我的要求 2,

我想忽略那些新添加的列并继续使用我之前的架构,即第 1 天的架构数据。

我的第 3 项要求,

如果传入的 csv 数据(即第 10 天架构)的架构发生变化,我还想自动添加新架构

【问题讨论】:

  • 你需要找到模式之间的差异,here你可以找到关于它的广泛讨论

标签: apache-spark pyspark apache-spark-sql spark-csv


【解决方案1】:
import org.apache.spark.sql.DataFrame

object SchemaDiff {

  def main(args: Array[String]): Unit = {
    // Just because its a simple CSV not considering column data type changes
    val df1 : DataFrame = null // Dataframe for yesterday's data
    val df2 : DataFrame = null // Dataframe for today's data
    val deltaColumnNames = df2.columns.diff(df1.columns)
    val ignoreSchemaChange = true
    if(!deltaColumnNames.isEmpty) {
      println("Schema change")
    }
    val resultDf = if(ignoreSchemaChange) {
      df2.toDF(df1.columns: _*) // Maintain yesterday's schema
    } else {
      df2  // Use updated schema
    }

  }
}

【讨论】:

    猜你喜欢
    • 2020-09-17
    • 2021-12-13
    • 2023-03-25
    • 1970-01-01
    • 2020-01-28
    • 2015-02-12
    • 1970-01-01
    • 2018-10-08
    • 1970-01-01
    相关资源
    最近更新 更多