【问题标题】:Spark scala aggregate to an array and concat itSpark scala 聚合到一个数组并连接它
【发布时间】:2022-01-24 02:48:27
【问题描述】:

我有一个包含许多列的数据集,如下所示:(Columns -name, timestamp, platform, clickcount, id)

Joy  2021-10-10T10:27:16  apple      5   1
May  2020-12-12T22:28:08  android    6   2
June 2021-09-15T20:20:06  Microsoft  9   3
Joy  2021-09-09T09:30:09  android    10  1
May  2021-08-08T05:05:05  apple      8   2

我想按 id 分组,之后应该是这样的

Joy  2021-10-10T10:27:16,2021-09-09T09:30:09   apple,android         5,10   1
May  2020-12-12T22:28:08,2021-08-08T05:05:05   android,apple         6,8    2
June 2021-09-15T20:20:06                       Microsoft             9      3

在调用另一个将 id 转换为伪 id 的 Api 后,我想映射该 id 并使其看起来像

Joy  2021-10-10T10:27:16,2021-09-09T09:30:09   apple,android         5,10   1   A12
May  2020-12-12T22:28:08,2021-08-08T05:05:05   android,apple         6,8    2   B23
June 2021-09-15T20:20:06                       Microsoft             9      3   C34

我已尝试使用 groupByforEach,但我被卡住了,无法继续进行

【问题讨论】:

    标签: scala apache-spark aggregate-functions


    【解决方案1】:

    为了应用您想要的聚合,您应该使用collect_set 作为聚合函数,并使用concat_ws 以便用逗号加入创建的数组:

    import org.apache.spark.sql.functions.{collect_set, concat_ws}
    import spark.implicits._
    
    val df: DataFrame = Seq(
      ("joy", "2021-10-10T10:27:16", "apple", 5, 1),
      ("may", "2020-12-12T22:28:08", "android", 6, 2),
      ("june", "2021-09-15T20:20:06", "microsoft", 9, 3),
      ("joy", "2021-09-09T09:30:09", "android", 10, 1),
      ("may", "2021-08-08T05:05:05", "apple", 8, 2)
    ).toDF("name", "timestamp", "platform", "clickcount", "id")
    
    df
      .groupBy("id")
      .agg(
        concat_ws(",", collect_set("timestamp")).as("timestamp"),
        concat_ws(",", collect_set("name")).as("name"),
        concat_ws(",", collect_set("platform")).as("platform"),
        concat_ws(",", collect_set("clickcount")).as("clickcount")
      ).show()
    

    输出应该是:

    +---+--------------------+----+-------------+----------+
    | id|           timestamp|name|     platform|clickcount|
    +---+--------------------+----+-------------+----------+
    |  1|2021-10-10T10:27:...| joy|apple,android|      5,10|
    |  3| 2021-09-15T20:20:06|june|    microsoft|         9|
    |  2|2021-08-08T05:05:...| may|apple,android|       6,8|
    +---+--------------------+----+-------------+----------+
    

    为了添加一个伪 id 列,您应该将创建的df 与另一个包含转换值的数据框连接起来,或者编写一个将接收 id 值并将其转换为伪 id 的 UDF。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-04-07
      • 1970-01-01
      • 2019-06-02
      • 2020-01-20
      • 1970-01-01
      相关资源
      最近更新 更多