【问题标题】:Merge Dataframes With Differents Schemas - Scala Spark合并具有不同模式的数据帧 - Scala Spark
【发布时间】:2020-08-27 14:14:44
【问题描述】:

我正在将 JSON 转换为数据框。在第一步中,我创建一个数据框数组,然后创建一个联合。但是我在使用不同模式的 JSON 中进行联合时遇到了问题。

如果 JSON 具有与您在另一个问题中看到的相同的架构,我可以做到:Parse JSON root in a column using Spark-Scala

我正在处理以下数据:

val exampleJsonDifferentSchema = spark.createDataset(

      """
      {"ITEM1512":
            {"name":"Yin",
             "address":{"city":"Columbus",
                        "state":"Ohio"},
             "age":28           }, 
        "ITEM1518":
            {"name":"Yang",
             "address":{"city":"Working",
                        "state":"Marc"}
                        },
        "ITEM1458":
            {"name":"Yossup",
             "address":{"city":"Macoss",
                        "state":"Microsoft"},
            "age":28
                        }
      }""" :: Nil)

如您所见,不同之处在于一个数据框没有年龄。

val itemsExampleDiff = spark.read.json(exampleJsonDifferentSchema)
itemsExampleDiff.show(false)
itemsExampleDiff.printSchema

+---------------------------------+---------------------------+-----------------------+
|ITEM1458                         |ITEM1512                   |ITEM1518               |
+---------------------------------+---------------------------+-----------------------+
|[[Macoss, Microsoft], 28, Yossup]|[[Columbus, Ohio], 28, Yin]|[[Working, Marc], Yang]|
+---------------------------------+---------------------------+-----------------------+

root
 |-- ITEM1458: struct (nullable = true)
 |    |-- address: struct (nullable = true)
 |    |    |-- city: string (nullable = true)
 |    |    |-- state: string (nullable = true)
 |    |-- age: long (nullable = true)
 |    |-- name: string (nullable = true)
 |-- ITEM1512: struct (nullable = true)
 |    |-- address: struct (nullable = true)
 |    |    |-- city: string (nullable = true)
 |    |    |-- state: string (nullable = true)
 |    |-- age: long (nullable = true)
 |    |-- name: string (nullable = true)
 |-- ITEM1518: struct (nullable = true)
 |    |-- address: struct (nullable = true)
 |    |    |-- city: string (nullable = true)
 |    |    |-- state: string (nullable = true)
 |    |-- name: string (nullable = true)

我现在的解决方案是如下代码,我在其中创建了一个 DataFrame 数组:

val columns:Array[String]       = itemsExample.columns
var arrayOfExampleDFs:Array[DataFrame] = Array()

for(col_name <- columns){

  val temp = itemsExample.select(lit(col_name).as("Item"), col(col_name).as("Value"))

  arrayOfExampleDFs = arrayOfExampleDFs :+ temp
}

val jsonDF = arrayOfExampleDFs.reduce(_ union _)

但是当我在联合中减少时,我有一个带有 不同架构 的 JSON,我不能这样做,因为数据框需要具有相同的架构。事实上,我有以下错误:

org.apache.spark.sql.AnalysisException: 联合只能在上执行 具有兼容列类型的表。

我正在尝试做一些我在这个问题中发现的类似事情:How to perform union on two DataFrames with different amounts of columns in spark?

具体那部分:

val cols1 = df1.columns.toSet
val cols2 = df2.columns.toSet
val total = cols1 ++ cols2 // union

def expr(myCols: Set[String], allCols: Set[String]) = {
  allCols.toList.map(x => x match {
    case x if myCols.contains(x) => col(x)
    case _ => lit(null).as(x)
  })
}

但我无法为列设置集合,因为我需要动态捕获列的总数和单项。我只能这样做:

for(i <- 0 until arrayOfExampleDFs.length-1) {

    val cols1 = arrayOfExampleDFs(i).select("Value").columns.toSet
    val cols2 = arrayOfExampleDFs(i+1).select("Value").columns.toSet
    val total = cols1 ++ cols2

    arrayOfExampleDFs(i).select("Value").printSchema()

    print(total)
}

那么,怎么可能是一个动态执行联合的函数呢?

更新:预期输出

在这种情况下,这个数据框和架构:

+--------+---------------------------------+
|Item    |Value                            |
+--------+---------------------------------+
|ITEM1458|[[Macoss, Microsoft], 28, Yossup]|
|ITEM1512|[[Columbus, Ohio], 28, Yin]      |
|ITEM1518|[[Working, Marc], null, Yang]    |
+--------+---------------------------------+

root
 |-- Item: string (nullable = false)
 |-- Value: struct (nullable = true)
 |    |-- address: struct (nullable = true)
 |    |    |-- city: string (nullable = true)
 |    |    |-- state: string (nullable = true)
 |    |-- age: long (nullable = true)
 |    |-- name: string (nullable = true)

【问题讨论】:

  • 如果您知道要考虑的所有可能的列,您可以简单地使用这些列创建一个列表并使用它。这样你就不需要在每次迭代中计算一个新的集合了。
  • 您还可以使用预定义架构读取 JSON,其中包含所有预期列,并且某些列也可以为空。如果您事先知道可能的模式的联合,则此方法有效。
  • 但在读取为 JSON 之前,我需要进行转换,因为项目(ITEM1458、ITEM1512、ITEM1518 等)显示为列,我需要将此列设为值。我可以在这里解决什么问题(对于具有相同架构的 json):stackoverflow.com/questions/61669258/…
  • @jqc 这个问题你解决了吗?
  • 不,我解决不了。 @AlexandrosBiratsis

标签: json scala apache-spark apache-spark-sql


【解决方案1】:

这是一种可能的解决方案,它通过在未找到年龄列时为所有数据帧创建通用架构:

import org.apache.spark.sql.functions.{col, lit, struct}
import org.apache.spark.sql.types.{LongType, StructField, StructType}

....

for(col_name <- columns){
  val currentDf = itemsExampleDiff.select(col(col_name))

  // try to identify if age field is present
  val hasAge = currentDf.schema.fields(0)
                        .dataType
                        .asInstanceOf[StructType]
                        .fields
                        .contains(StructField("age", LongType, true))

  val valueCol = hasAge match {
    // if not construct a new value column
    case false => struct(
                    col(s"${col_name}.address"), 
                    lit(null).cast("bigint").as("age"),
                    col(s"${col_name}.name")
                  )

    case true => col(col_name)
  }

  arrayOfExampleDFs = arrayOfExampleDFs :+ currentDf.select(lit(col_name).as("Item"), valueCol.as("Value"))
}

val jsonDF = arrayOfExampleDFs.reduce(_ union _)

// +--------+---------------------------------+
// |Item    |Value                            |
// +--------+---------------------------------+
// |ITEM1458|[[Macoss, Microsoft], 28, Yossup]|
// |ITEM1512|[[Columbus, Ohio], 28, Yin]      |
// |ITEM1518|[[Working, Marc],, Yang]         |
// +--------+---------------------------------+

分析:可能最苛刻的部分是找出age 是否存在。对于查找,我们使用df.schema.fields 属性,它允许我们深入研究每一列的内部模式。

当没有找到年龄时,我们使用struct 重新生成列:

struct(
   col(s"${col_name}.address"), 
   lit(null).cast("bigint").as("age"),
   col(s"${col_name}.name")
)

【讨论】:

    猜你喜欢
    • 2020-06-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-31
    • 2018-01-06
    • 1970-01-01
    • 1970-01-01
    • 2018-03-10
    相关资源
    最近更新 更多