【问题标题】:How to Transform a Spark Scala Nested Map within a Map Data Structure?如何在地图数据结构中转换 Spark Scala 嵌套地图?
【发布时间】:2020-04-23 01:18:02
【问题描述】:

我想使用 Scala 案例类的数组编写一个嵌套数据结构,该结构由另一个 Map 中的 Map 组成。

结果应该转换这个数据框:

|Value|Country| Timestamp| Sum|
+-----+-------+----------+----+
|  123|    ITA|1475600500|18.0|
|  123|    ITA|1475600516|19.0|
+-----+-------+----------+----+

进入:

+--------------------------------------------------------------------+
|value                                                               |
+--------------------------------------------------------------------+
[{"value":123,"attributes":{"ITA":{"1475600500":18,"1475600516":19}}}]
+--------------------------------------------------------------------+

下面的actualResult 数据集让我很接近,但结构与我预期的数据框不太一样。

case class Record(value: Integer, attributes: Map[String, Map[String, BigDecimal]])
val actualResult = df
  .map(r =>
    Array(
      Record(
        r.getAs[Int]("Value"),
        Map(
          r.getAs[String]("Country") ->
            Map(
              r.getAs[String]("Timestamp") -> new BigDecimal(
                r.getAs[Double]("Sum").toString
              )
            )
        )
      )
    )
  )

actualResult 数据集中的 Timestamp 列不会合并到同一 Record 行,而是创建两个单独的行。

+----------------------------------------------------+
|value                                               |
+----------------------------------------------------+
[{"value":123,"attributes":{"ITA":{"1475600516":19}}}]
[{"value":123,"attributes":{"ITA":{"1475600500":18}}}]
+----------------------------------------------------+

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    使用groupBycollect_list by creatng 使用结构的组合列,我能够获得如下输出的单行。

    val mycsv =
        """
          |Value|Country|Timestamp|Sum
          |  123|ITA|1475600500|18.0
          |  123|ITA|1475600516|19.0
        """.stripMargin('|').lines.toList.toDS()
    
    
      val df: DataFrame = spark.read.option("header", true)
        .option("sep", "|")
        .option("inferSchema", true)
        .csv(mycsv)
      df.show
    
      val df1 = df.
        groupBy("Value","Country")
        .agg(  collect_list(struct(col("Country"), col("Timestamp"), col("Sum"))).alias("attributes")).drop("Country")
    
    
      val json = df1.toJSON // you can save in to file
      json.show(false)
    

    结果合并 2 行

    +-----+-------+----------+----+
    |Value|Country| Timestamp| Sum|
    +-----+-------+----------+----+
    |123.0|ITA    |1475600500|18.0|
    |123.0|ITA    |1475600516|19.0|
    +-----+-------+----------+----+
    
    +----------------------------------------------------------------------------------------------------------------------------------------------+
    |value                                                                                                                                         |
    +----------------------------------------------------------------------------------------------------------------------------------------------+
    |{"Value":123.0,"attributes":[{"Country":"ITA","Timestamp":1475600500,"Sum":18.0},{"Country":"ITA","Timestamp":1475600516,"Sum":19.0}]}|
    +----------------------------------------------------------------------------------------------------------------------------------------------+
    
    

    【讨论】:

    • 更新了我的答案和一些列标签如何进入 json 对你来说可以吗?
    • 你能检查一下结果吗
    • 嗨,Eric,让我知道您对此的反馈。如有任何问题,请提出。
    猜你喜欢
    • 2018-03-15
    • 2020-08-12
    • 2020-12-14
    • 1970-01-01
    • 2020-06-30
    • 1970-01-01
    • 1970-01-01
    • 2016-11-14
    • 1970-01-01
    相关资源
    最近更新 更多