【问题标题】:Scala Spark GroupBy Aggregate maintain order of dates while sorting listScala Spark GroupBy Aggregate 在排序列表时维护日期顺序
【发布时间】:2019-01-20 16:02:27
【问题描述】:
val df = Seq((1221, 1, "Boston", "9/22/18 14:00"), (1331, 1, "New York", "8/10/18 14:00"), (1442, 1, "Toronto", "10/15/19 14:00"), (2041, 2, "LA", "1/2/18 14:00"), (2001, 2,"San Fransisco", "5/20/18 15:00"), (3001, 3, "San Jose", "6/02/18 14:00"), (3121, 3, "Seattle", "9/12/18 16:00"), (34562, 3, "Utah", "12/12/18 14:00"), (3233, 3, "Boston", "8/31/18 14:00"), (4120, 4, "Miami", "1/01/18 14:00"), (4102, 4, "Cincinati", "7/21/19 14:00"), (4201, 4, "Washington", "5/10/18 23:00"), (4301, 4, "New Jersey", "3/27/18 15:00"), (4401, 4, "Raleigh", "11/14/18 14:00")).toDF("id", "group_id", "place", "date")

这是一个简单的df

|   id|group_id|        place|          date|  
+-----+--------+-------------+--------------+  
| 1221|       1|       Boston| 9/22/18 14:00|  
| 1331|       1|     New York| 8/10/18 14:00|    
| 1442|       1|      Toronto|10/15/19 14:00|  
| 2041|       2|           LA|  1/2/18 14:00|  
| 2001|       2|San Fransisco| 5/20/18 15:00|  
| 3001|       3|     San Jose| 6/02/18 14:00|  
| 3121|       3|      Seattle| 9/12/18 16:00|  
| 4562|       3|         Utah|12/12/18 14:00|  
| 3233|       3|       Boston| 8/31/18 14:00|  
| 4120|       4|        Miami| 1/01/18 14:00|  
| 4102|       4|    Cincinati| 7/21/19 14:00|  
| 4201|       4|   Washington| 5/10/18 23:00|  
| 4301|       4|   New Jersey| 3/27/18 15:00|  
| 4401|       4|      Raleigh|11/14/18 14:00|  
+-----+--------+-------------+--------------+  

我想按“group_id”分组并按升序收集日期。 (最早的日期在前)。

需要的输出:

+--------+----+--------+--------------+----+-------------+--------------+----+----------+--------------+----+---------+--------------+
|group_id|id_1| venue_1|        date_1|id_2|      venue_2|        date_2|id_3|   venue_3|        date_3|id_4|  venue_4|        date_4|
+--------+----+--------+--------------+----+-------------+--------------+----+----------+--------------+----+---------+--------------+
|       1|1331|New York|08/10/18 14:00|1221|       Boston|09/22/18 14:00|1442|   Toronto|10/15/19 14:00|null|     null|          null|
|       3|3001|San Jose|06/02/18 14:00|3233|       Boston|08/31/18 14:00|3121|   Seattle|09/12/18 16:00|4562|     Utah|12/12/18 14:00|
|       4|4120|   Miami|01/01/18 14:00|4301|   New Jersey|03/27/18 15:00|4201|Washington|05/10/18 23:00|4102|Cincinati|07/21/19 14:00|
|       2|2041|      LA| 01/2/18 14:00|2001|San Fransisco|05/20/18 15:00|null|      null|          null|null|     null|          null|
+--------+----+--------+--------------+----+-------------+--------------+----+----------+--------------+----+---------+--------------+

我正在使用的代码:

//for sorting by date to preserve order
val df2 = df.repartition(col("group_id")).sortWithinPartitions("date")

val finalDF = df2.groupBy(df("group_id")).agg(collect_list(df("id")).alias("id_list"),collect_list(df("place")).alias("venue_name_list"),collect_list(df("date")).alias("date_list")).selectExpr("group_id","id_list[0] as id_1","venue_name_list[0] as venue_1","date_list[0] as date_1","id_list[1] as id_2","venue_name_list[1] as venue_2","date_list[1] as date_2","id_list[2] as id_3","venue_name_list[2] as venue_3","date_list[2] as date_3","id_list[3] as id_4","venue_name_list[3] as venue_4","date_list[3] as date_4")

但输出是:

+--------+-----+-------+--------------+----+-------------+--------------+----+----------+-------------+----+----------+-------------+
|group_id| id_1|venue_1|        date_1|id_2|      venue_2|        date_2|id_3|   venue_3|       date_3|id_4|   venue_4|       date_4|
+--------+-----+-------+--------------+----+-------------+--------------+----+----------+-------------+----+----------+-------------+
|       1| 1442|Toronto|10/15/19 14:00|1331|     New York| 8/10/18 14:00|1221|    Boston|9/22/18 14:00|null|      null|         null|
|       3|34562|   Utah|12/12/18 14:00|3001|     San Jose| 6/02/18 14:00|3233|    Boston|8/31/18 14:00|3121|   Seattle|9/12/18 16:00|
|       4| 4120|  Miami| 1/01/18 14:00|4401|      Raleigh|11/14/18 14:00|4301|New Jersey|3/27/18 15:00|4201|Washington|5/10/18 23:00|
|       2| 2041|     LA|  1/2/18 14:00|2001|San Fransisco| 5/20/18 15:00|null|      null|         null|null|      null|         null|
+--------+-----+-------+--------------+----+-------------+--------------+----+----------+-------------+----+----------+-------------+

观察: 如果日期被格式化而不是例如“9/22/18 14:00”到“09/22/18 14:00”,则在一位数月份日期之前添加“0”并在一位数日期之前添加零,代码工作正常,也就是说,日期顺序得到了正确维护。欢迎任何解决方案!谢谢你。

【问题讨论】:

    标签: scala apache-spark dataframe group-by aggregate


    【解决方案1】:

    正如您已经发现按未格式化的StringType 日期排序是问题的根源,这是一种首先生成TimestampType 日期的方法,在“合适的字段顺序”中创建所需列的StructType 列" 用于排序:

    val finalDF = df.
      withColumn("dateFormatted", to_timestamp($"date", "MM/dd/yy HH:mm")).
      groupBy($"group_id").agg(
        sort_array(collect_list(struct($"dateFormatted", $"id", $"place"))).as("sorted_arr")
      ).
      selectExpr(
        "group_id",
        "sorted_arr[0].id as id_1", "sorted_arr[0].place as venue_1", "sorted_arr[0].dateFormatted as date_1",
        "sorted_arr[1].id as id_2", "sorted_arr[1].place as venue_2", "sorted_arr[1].dateFormatted as date_2",
        "sorted_arr[2].id as id_3", "sorted_arr[2].place as venue_3", "sorted_arr[2].dateFormatted as date_3",
        "sorted_arr[3].id as id_4", "sorted_arr[3].place as venue_4", "sorted_arr[3].dateFormatted as date_4"
      )
    
    finalDF.show
    // +--------+----+--------+-------------------+----+-------------+-------------------+----+----------+-------------------+-----+-------+-------------------+
    // |group_id|id_1| venue_1|             date_1|id_2|      venue_2|             date_2|id_3|   venue_3|             date_3| id_4|venue_4|             date_4|
    // +--------+----+--------+-------------------+----+-------------+-------------------+----+----------+-------------------+-----+-------+-------------------+
    // |       1|1331|New York|2018-08-10 14:00:00|1221|       Boston|2018-09-22 14:00:00|1442|   Toronto|2019-10-15 14:00:00| null|   null|               null|
    // |       3|3001|San Jose|2018-06-02 14:00:00|3233|       Boston|2018-08-31 14:00:00|3121|   Seattle|2018-09-12 16:00:00|34562|   Utah|2018-12-12 14:00:00|
    // |       4|4120|   Miami|2018-01-01 14:00:00|4301|   New Jersey|2018-03-27 15:00:00|4201|Washington|2018-05-10 23:00:00| 4401|Raleigh|2018-11-14 14:00:00|
    // |       2|2041|      LA|2018-01-02 14:00:00|2001|San Fransisco|2018-05-20 15:00:00|null|      null|               null| null|   null|               null|
    // +--------+----+--------+-------------------+----+-------------+-------------------+----+----------+-------------------+-----+-------+-------------------+
    

    几点说明:

    1. 必须形成StructType 列以确保将相应的列排序在一起
    2. 结构字段dateFormatted 放在第一位,这样sort_array 将按所需顺序对数组进行排序

    【讨论】:

    • 分区和排序并不总是在维护的顺序中给出相同的结果。即使将其转换为 TimeStamptype。
    • @saurin shah,请参阅修改后的答案以解决上述问题。
    • 这确实有帮助。但是我仍然发现对于较大的数据集没有维护顺序的情况。 SPARK 2.1 版本可能不支持它。
    • @saurin shah,我会检查不按顺序排列的行(特别是 dataFormatted 的值)。如果原始日期字符串并非全部采用“MM/dd/yy HH:mm”格式,则某些dataFormatted 将不正确并且排序将被关闭。
    • 是的。可能这就是订单关闭的原因。我将检查日期的字符串格式。再次感谢
    【解决方案2】:

    使用 to_timestamp 函数格式化日期并使用 sort_array 在聚合中排序,如下所示:

     import org.apache.spark.sql.functions.to_timestamp
    
      val df = Seq((1221, 1, "Boston", "9/22/18 14:00"), (1331, 1, "New York", "8/10/18 14:00"), (1442, 1, "Toronto", "10/15/19 14:00"), (2041, 2, "LA", "1/2/18 14:00"), (2001, 2, "San Fransisco", "5/20/18 15:00"), (3001, 3, "San Jose", "6/02/18 14:00"), (3121, 3, "Seattle", "9/12/18 16:00"), (34562, 3, "Utah", "12/12/18 14:00"), (3233, 3, "Boston", "8/31/18 14:00"), (4120, 4, "Miami", "1/01/18 14:00"), (4102, 4, "Cincinati", "7/21/19 14:00"), (4201, 4, "Washington", "5/10/18 23:00"), (4301, 4, "New Jersey", "3/27/18 15:00"), (4401, 4, "Raleigh", "11/14/18 14:00"))
        .toDF("id", "group_id", "place", "date")
    
      val df2 = df.withColumn("MyDate", to_timestamp($"date", "MM/dd/yyyy HH:mm"))
    
      val finalDF = df2.groupBy(df("group_id"))
        .agg(collect_list(df2("id")).alias("id_list"),
          collect_list(df2("place")).alias("venue_name_list"),
          sort_array(collect_list(df2("MyDate"))).alias("date_list")).
        selectExpr("group_id",
          "id_list[0] as id_1",
          "venue_name_list[0] as venue_1",
          "date_list[0] as date_1",
          "id_list[1] as id_2",
          "venue_name_list[1] as venue_2",
          "date_list[1] as date_2",
          "id_list[2] as id_3",
          "venue_name_list[2] as venue_3",
          "date_list[2] as date_3",
          "id_list[3] as id_4",
          "venue_name_list[3] as venue_4",
          "date_list[3] as date_4")
    
    
      finalDF.show()
    

    【讨论】:

    • 太棒了。但是对于更大的数据集,订单维护。分区和排序并不总是以维持的顺序给出相同的结果。大概这种方法可以。但即使这个没有。谢谢。
    猜你喜欢
    • 1970-01-01
    • 2020-07-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多