【问题标题】:Transform structured data to JSON format using Spark Scala使用 Spark Scala 将结构化数据转换为 JSON 格式
【发布时间】:2020-01-21 02:45:46
【问题描述】:

我的“结构化数据”如下所示,我需要将其转换为下面显示的“预期结果”类型。我的“输出模式”也显示出来了。感谢您是否可以就我如何使用 Spark Scala 代码实现这一点提供一些帮助。

注意:对结构化数据进行分组是在 SN 和 VIN 列上完成的。 相同的SN 和VIN 应该有一行,如果SN 或VIN 发生变化,则数据将出现在下一行。

结构化数据:

+-----------------+-------------+--------------------+---+
|VIN              |ST           |SV                  |SN |
|FU74HZ501740XXXXX|1566799999225|44.0                |APP|
|FU74HZ501740XXXXX|1566800002758|61.0                |APP|
|FU74HZ501740XXXXX|1566800009446|23.39               |ASP|

预期结果:

输出架构:

val outputSchema = StructType(
  List(
    StructField("VIN", StringType, true),
    StructField("EVENTS", ArrayType(
        StructType(Array(
          StructField("SN", StringType, true),
          StructField("ST", IntegerType, true),
          StructField("SV", DoubleType, true)
        ))))
  )
)

【问题讨论】:

  • 请在您的问题中添加文本而不是图像。它使我们更容易根据您当前的数据集重现问题。
  • 您在这里按哪一列分组? SN?
  • 你好@Shaido,是的,应该在列 SN 上进行分组...我想遍历 SN 列,并将 SN、ST、SV 包括在列 EVENTS 和 VIN 的单个数组中另一列如预期结果所示。
  • @AnilKumarKB:如果 SN 对同一个 VIN 有不同的 valeus 会发生什么?例如,在您的示例中,如果第二行中的 VIN 与第一行不同。
  • 您好@Shaido 很抱歉造成混淆,它应该按 SN 和 VIN 分组。例如:一行代表相同的 SN 和 VIN。如果 SN 或 VIN 发生变化,则数据将出现在下一行中。

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


【解决方案1】:

您可以通过 SparkSession 获得它。


val df = spark.read.json("/path/to/json/file/test.json")

这里的 spark 是 SparkSession 对象

【讨论】:

  • 请注意,问题不是关于读取 json 文件,而是将数据框转换为具有包含 json 元素的数组列。
  • 我理解可以转换成json的结构化数据,没有把它当成dataframe。我的错。
【解决方案2】:

从 Spark 2.1 开始,您可以使用 struct 和 collect_list 实现此目的。

val df_2 = Seq(
  ("FU74HZ501740XXXX",1566799999225.0,44.0,"APP"),
  ("FU74HZ501740XXXX",1566800002758.0,61.0,"APP"),
  ("FU74HZ501740XXXX",1566800009446.0,23.39,"ASP")
).toDF("vin","st","sv","sn") 

df_2.show(false)
+----------------+-----------------+-----+---+
|vin             |st               |sv   |sn |
+----------------+-----------------+-----+---+
|FU74HZ501740XXXX|1.566799999225E12|44.0 |APP|
|FU74HZ501740XXXX|1.566800002758E12|61.0 |APP|
|FU74HZ501740XXXX|1.566800009446E12|23.39|ASP|
+----------------+-----------------+-----+---+

使用collect_list 和struct:

df_2.groupBy("vin","sn")
  .agg(collect_list(struct($"st", $"sv",$"sn")).as("events"))
  .withColumn("events",to_json($"events"))
  .drop(col("sn"))

这将给出想要的输出:

+----------------+---------------------------------------------------------------------------------------------+
|vin             |events                                                                                       |
+----------------+---------------------------------------------------------------------------------------------+
|FU74HZ501740XXXX|[{"st":1.566800009446E12,"sv":23.39,"sn":"ASP"}]                                             |
|FU74HZ501740XXXX|[{"st":1.566799999225E12,"sv":44.0,"sn":"APP"},{"st":1.566800002758E12,"sv":61.0,"sn":"APP"}]|
+----------------+---------------------------------------------------------------------------------------------+

【讨论】:

  • +1。您可以在agg (collect_list(to_json(struct(...) 中添加to_json,以避免出现额外的withColumn。旧版本的 Spark 也应该可以做到这一点,可能是在引入 to_json 时的 2.1 版。
  • @AnilKumarKB 很棒
猜你喜欢
  • 1970-01-01
  • 2021-10-05
  • 1970-01-01
  • 2020-01-08
  • 1970-01-01
  • 2018-03-28
  • 2015-01-05
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多