【发布时间】:2020-09-03 09:57:31
【问题描述】:
我有一个代码 sn-p,它将读取文件路径的 Json 数组,然后合并输出并给我两个不同的表。所以我想为这两个表创建两个不同的 createOrReplaceview(name) 并且名称将在 json 数组中可用,如下所示:
{
"source": [
{
"name": "testPersons",
"data": [
"E:\\dataset\\2020-05-01\\",
"E:\\dataset\\2020-05-02\\"
],
"type": "json"
},
{
"name": "testPets",
"data": [
"E:\\dataset\\2020-05-01\\078\\",
"E:\\dataset\\2020-05-02\\078\\"
],
"type": "json"
}
]
}
我的输出:
testPersons
+---+------+
|name |age|
+---+------+
|John |24 |
|Cammy |20 |
|Britto|30 |
|George|23 |
|Mikle |15 |
+---+------+
testPets
+---+------+
|name |age|
+---+------+
|piku |2 |
|jimmy |3 |
|rapido|1 |
+---+------+
上面是我的输出和 Json 数组我的代码遍历每个数组并读取数据部分并读取数据。
但是如何更改我的以下代码为每个输出表创建一个临时视图。
例如我想创建.createOrReplaceTempView(testPersons) 和.createOrReplaceTempView(testPets)
根据 Json 数组查看名称
if (dataArr(counter)("type").value.toString() == "json") {
val name = dataArr(counter)("name").value.toString()
val dataPath = dataArr(counter)("data").arr
val input = dataPath.map(item => {
val rdd = spark.sparkContext.wholeTextFiles(item.str).map(i => "[" + i._2.replaceAll("\\}.*\n{0,}.*\\{", "},{") + "]")
spark
.read
.schema(Schema.getSchema(name))
.option("multiLine", true)
.json(rdd)
})
val emptyDF = spark.createDataFrame(spark.sparkContext.emptyRDD[Row], Schema.getSchema(name))
val finalDF = input.foldLeft(emptyDF)((x, y) => x.union(y))
finalDF.show()
预期输出:
spark.sql("SELECT * FROM testPersons").show()
spark.sql("SELECT * FROM testPets").show()
它应该给我一张桌子。
【问题讨论】:
标签: scala apache-spark