【问题标题】:Scala - How to convert JSON Keys and Values as columnsScala - 如何将 JSON 键和值转换为列
【发布时间】:2020-07-11 14:50:32
【问题描述】:

如何将下面的 Input Json 解析为 key 和 value 列。任何帮助表示赞赏。

输入:

{
"name" : "srini",
"value": {
"1" : "val1",
"2" : "val2",
"3" : "val3"
}
}

    Output DataFrame Column:

    name      key        value
    -----------------------------
    srini      1         val1
    srini      2         val2
    srini      3         val3



        //++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++Input DataFrame :
        +--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
        |json_file                                                                                                                                                                                                                                                                                                     |
        +--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
        |{"file_path":"AAA/BBB.CCC.zip","file_name":"AAA_20200202122754.json","received_time":"2020-03-31","obj_cls":"Monitor","obj_cls_inst":"Monitor","relation_tree":"Source~>HD_Info~>Monitor","s_tag":"ABC1234","Monitor":{"Index":"0","Vendor_Data":"58F5Y","Monitor_Type":"Lenovo Monitor","HnfoID":"650FEC74"}}| 
        +--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+


        How to convert this above json file in a DataFrame like below :

        +----------------+-----------------------+--------------+--------+-------------+-------------------------+----------+----------------+----------------+
        |file_path       |file_name              |received_time |obj_cls |obj_cls_inst |relation_tree            |s_tag     |attribute_name  |attribute_value |
        +----------------+-----------------------+--------------+--------+-------------+-------------------------+----------+----------------+----------------+
        |AAA/BBB.CCC.zip |AAA_20200202122754.json|2020-03-31    |Monitor |Monitor      |Source~>HD_Info~>Monitor |ABC1234   |Index           |0               |
        +----------------+-----------------------+--------------+--------+-------------+-------------------------+----------+----------------+----------------+
        |AAA/BBB.CCC.zip |AAA_20200202122754.json|2020-03-31    |Monitor |Monitor      |Source~>HD_Info~>Monitor |ABC1234   |Vendor_Data     |58F5Y           |
        +----------------+-----------------------+--------------+--------+-------------+-------------------------+----------+----------------+----------------+
        |AAA/BBB.CCC.zip |AAA_20200202122754.json|2020-03-31    |Monitor |Monitor      |Source~>HD_Info~>Monitor |ABC1234   |Monitor_Type    |Lenovo Monitor  |
        +----------------+-----------------------+--------------+--------+-------------+-------------------------+----------+----------------+----------------+
        |AAA/BBB.CCC.zip |AAA_20200202122754.json|2020-03-31    |Monitor |Monitor      |Source~>HD_Info~>Monitor |ABC1234   |HnfoID          |650FEC74        |
        +----------------+-----------------------+--------------+--------+-------------+-------------------------+----------+----------------+----------------+

//**********************************************
val rawData = sparkSession.sql("select 1").withColumn("obj_cls", lit("First")).withColumn("s_tag", lit("S_12345")).withColumn("jsonString", lit("""{"id":""1,"First":{"Info":"ABCD123","Res":"5.2"}}"""))

【问题讨论】:

  • 嗨@SCouto,非常感谢您的回复,非常有帮助。实际上我已经添加了有问题的输入数据帧和预期输出数据帧。如果可能,请查看并提供帮助:(
  • 我用新的例子更新了我的提示,但它是一样的,你只需要更改列名。请测试一下,如果没问题,请接受我的问题,如果其他人有同样的问题,可以快速找到答案

标签: json scala apache-spark parsing key


【解决方案1】:

将 json 加载到 DF 中后,如下所示:

+-----+------------------+
| name|             value|
+-----+------------------+
|srini|[val1, val2, val3]|
+-----+------------------+

首先选择整个值项:

df.select($"name", $"value.*")

这会给你这个:

+-----+----+----+----+
| name|   1|   2|   3|
+-----+----+----+----+
|srini|val1|val2|val3|
+-----+----+----+----+

然后你需要将列转成行,为此我通常定义一个辅助函数kv:

 def kv (columnsToTranspose: Array[String]) = explode(array(columnsToTranspose.map {
    c => struct(lit(c).alias("k"), col(c).alias("v"))
  }: _*))

然后为所需的列创建一个数组:

val pivotCols = Array("1", "2", "3")

最后将函数应用到之前的DF:

df.select($"name", $"value.*")
.withColumn("kv", kv(pivotCols))
.select($"name", $"kv.k" as "key", $"kv.v" as "value")

结果:

+-----+---+-----+
| name|key|value|
+-----+---+-----+
|srini|  1| val1|
|srini|  2| val2|
|srini|  3| val3|
+-----+---+-----+

编辑

如果您不想手动指定要旋转的列,可以使用中间 df,如下所示:

val dfIntermediate = df.select($"name", $"value.*")

dfIntermediate.withColumn("kv", kv(dfIntermediate.columns.tail))
.select($"name", $"kv.k" as "key", $"kv.v" as "value")

你会得到同样的结果:

+-----+---+-----+
| name|key|value|
+-----+---+-----+
|srini|  1| val1|
|srini|  2| val2|
|srini|  3| val3|
+-----+---+-----+

EDIT2

与新示例相同,您只需要更改您读取/透视的列

val pivotColumns = Array("HnfoId", "Index", "Monitor_Type", "Vendor_Data")

df.select("file_path", "file_name", "received_time", "obj_cls", "obj_cls_inst", "relation_tree", "s_Tag", "Monitor.*").withColumn("kv", kv(pivotColumns)).select($"file_path", $"file_name", $"received_time", $"obj_cls", $"obj_cls_inst", $"relation_tree", $"s_Tag", $"kv.k" as "attribute_name", $"kv.v" as "attribute_value").show
+---------------+--------------------+-------------+-------+------------+--------------------+-------+--------------+---------------+
|      file_path|           file_name|received_time|obj_cls|obj_cls_inst|       relation_tree|  s_Tag|attribute_name|attribute_value|
+---------------+--------------------+-------------+-------+------------+--------------------+-------+--------------+---------------+
|AAA/BBB.CCC.zip|AAA_2020020212275...|   2020-03-31|Monitor|     Monitor|Source~>HD_Info~>...|ABC1234|        HnfoId|       650FEC74|
|AAA/BBB.CCC.zip|AAA_2020020212275...|   2020-03-31|Monitor|     Monitor|Source~>HD_Info~>...|ABC1234|         Index|              0|
|AAA/BBB.CCC.zip|AAA_2020020212275...|   2020-03-31|Monitor|     Monitor|Source~>HD_Info~>...|ABC1234|  Monitor_Type| Lenovo Monitor|
|AAA/BBB.CCC.zip|AAA_2020020212275...|   2020-03-31|Monitor|     Monitor|Source~>HD_Info~>...|ABC1234|   Vendor_Data|          58F5Y|
+---------------+--------------------+-------------+-------+------------+--------------------+-------+--------------+---------------+

【讨论】:

  • 嗨 SCouto,再次感谢,但问题是,我有 100 个带有动态字段的 Json 文件,所有的 json 文件都有不同的 Pivot 列,那么有没有什么办法可以在 pivotColumns 中不使用硬代码?
  • 您可以通过从结构中提取内部字段来完成(这样您只需要知道结构名称) df.select($"Monitor.*").schema.names
  • 嗨,@SCouto,你能帮忙解决最后一个问题吗? Json 的结构字段(例如:First)将始终是 obj_cls 的值,以便在不同 Json 数据的不同结构字段的情况下动态获取它。
猜你喜欢
  • 1970-01-01
  • 2020-02-05
  • 2022-08-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-01-06
  • 1970-01-01
相关资源
最近更新 更多