【发布时间】: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