【问题标题】:Flatten nested json in Scala Spark Dataframe在 Scala Spark Dataframe 中展平嵌套的 json
【发布时间】:2021-03-23 12:48:43
【问题描述】:

我有多个来自任何 restapi 的 json,但我不知道它的架构。我无法使用 dataframes 的 explode 功能,因为我不知道由 spark api 创建的列名。

1.我们能否通过解码dataframe.schema.fields中的值来存储嵌套数组元素键的键,因为spark只提供数据帧行中的值部分,并将顶级键作为列名。

数据框 --

+--------------------+
|       stackoverflow|
+--------------------+
|[[[Martin Odersky...|
+--------------------+

是否有任何最佳方式通过在运行时确定架构来使用数据框方法来展平 json。

示例 Json -:

{
  "stackoverflow": [{
    "tag": {
      "id": 1,
      "name": "scala",
      "author": "Martin Odersky",
      "frameworks": [
        {
          "id": 1,
          "name": "Play Framework"
        },
        {
          "id": 2,
          "name": "Akka Framework"
        }
      ]
    }
  },
    {
      "tag": {
        "id": 2,
        "name": "java",
        "author": "James Gosling",
        "frameworks": [
          {
            "id": 1,
            "name": "Apache Tomcat"
          },
          {
            "id": 2,
            "name": "Spring Boot"
          }
        ]
      }
    }
  ]
}

注意 - 我们需要在 dataframe 中进行所有操作,因为有大量数据即将到来,我们无法解析每一个 json。

【问题讨论】:

    标签: scala apache-spark multidimensional-array apache-spark-sql


    【解决方案1】:

    尽量避免将所有列展平。

    创建了辅助函数&你可以直接在DataFrame上调用df.explodeColumns

    以下代码将展平多级数组和结构类型列。

    scala> :paste
    // Entering paste mode (ctrl-D to finish)
    
    import org.apache.spark.sql.{DataFrame, SparkSession}
    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.types._
    import scala.annotation.tailrec
    import scala.util.Try
    
    implicit class DFHelpers(df: DataFrame) {
        def columns = {
          val dfColumns = df.columns.map(_.toLowerCase)
          df.schema.fields.flatMap { data =>
            data match {
              case column if column.dataType.isInstanceOf[StructType] => {
                column.dataType.asInstanceOf[StructType].fields.map { field =>
                  val columnName = column.name
                  val fieldName = field.name
                  col(s"${columnName}.${fieldName}").as(s"${columnName}_${fieldName}")
                }.toList
              }
              case column => List(col(s"${column.name}"))
            }
          }
        }
    
        def flatten: DataFrame = {
          val empty = df.schema.filter(_.dataType.isInstanceOf[StructType]).isEmpty
          empty match {
            case false =>
              df.select(columns: _*).flatten
            case _ => df
          }
        }
        def explodeColumns = {
          @tailrec
          def columns(cdf: DataFrame):DataFrame = cdf.schema.fields.filter(_.dataType.typeName == "array") match {
            case c if !c.isEmpty => columns(c.foldLeft(cdf)((dfa,field) => {
              dfa.withColumn(field.name,explode_outer(col(s"${field.name}"))).flatten
            }))
            case _ => cdf
          }
          columns(df.flatten)
        }
    }
    
    // Exiting paste mode, now interpreting.
    
    import org.apache.spark.sql.{DataFrame, SparkSession}
    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.types._
    import scala.annotation.tailrec
    import scala.util.Try
    defined class DFHelpers
    
    

    扁平列

    scala> df.printSchema
    root
     |-- stackoverflow: array (nullable = true)
     |    |-- element: struct (containsNull = true)
     |    |    |-- tag: struct (nullable = true)
     |    |    |    |-- author: string (nullable = true)
     |    |    |    |-- frameworks: array (nullable = true)
     |    |    |    |    |-- element: struct (containsNull = true)
     |    |    |    |    |    |-- id: long (nullable = true)
     |    |    |    |    |    |-- name: string (nullable = true)
     |    |    |    |-- id: long (nullable = true)
     |    |    |    |-- name: string (nullable = true)
    
    
    scala> df.explodeColumns.printSchema
    root
     |-- author: string (nullable = true)
     |-- frameworks_id: long (nullable = true)
     |-- frameworks_name: string (nullable = true)
     |-- id: long (nullable = true)
     |-- name: string (nullable = true)
    
    scala>
    
    

    【讨论】:

    • 您有什么问题吗?
    • 抱歉之前的评论。如果我在"npi": { "$numberInt": "2038094571" }, "part_names": "VALARIE", "salary": { "$numberInt": "250000" } } 上方给定数组的底部附加这个额外的部分' 是模棱两可的,可能是:$numberInt, $numberInt.;在 colums_* 方法中找不到子文档重复键。
    • @Srinivas- 我正在尝试使用 spark 结构化流来完成相同的场景。使用动态模式从kafka读取jsonstring(每条记录可以有不同的模式;一条记录可以有5列,其他记录可以有3列)。这是一个嵌套的,我将其展平并使用 foreachwriter 将其加载到 hbase。但我真正的问题是模式推断,没有模式。 Spark 不支持它。任何线索都会很有帮助。谢谢
    • 来自 Kafka 的单个名为 value 的列,其中包含 jsonstring 将通过流式传输读取。 Schema 不会被定义。它会倾向于改变。我们需要在加载到 hbase 之前解析和展平该 jsonstring。
    • @Srinivas 您好,您的方法令人印象深刻。只是好奇为什么 Spark 本身不支持这种功能(我认为这是一个非常常见的用例),你知道吗?
    猜你喜欢
    • 1970-01-01
    • 2020-10-02
    • 1970-01-01
    • 2019-09-26
    • 2016-12-01
    • 1970-01-01
    • 2016-07-19
    • 2018-08-13
    • 2018-10-28
    相关资源
    最近更新 更多