【问题标题】:collect_list in scala dataframe which will collect rows in an interval of fixed column numbersscala数据框中的collect_list,它将以固定列号的间隔收集行
【发布时间】:2019-11-29 05:15:11
【问题描述】:

我需要将特定分区的所有行收集到数据框中的一行中。 我必须将此数据帧转储到 cosmosDB 中,每个文档只能保存 2MB 的数据。 但是当我将上述数据帧收集到一行时,它超过了 2MB 并在写入 CosmosDB 时抛出错误。

我想将这些行合并为一个,固定间隔为 500 行。 对于一个分区,前 500 行应该被收集到一行中,接下来的 500 行应该被收集到另一行中,依此类推..

输入数据如下图。

+------+----------+---------------------------------------+
|ID    |TIME      |SGNL                                   |
+------+----------+---------------------------------------+
|00001 |1574360355|{"SN":"Acc","ST":1574360296,"SV":"0.0"}|
|00001 |1574360355|{"SN":"Acc","ST":1574360296,"SV":"0.0"}|
|00001 |1574360355|{"SN":"Acc","ST":1574360296,"SV":"0.0"}|
|00001 |1574360355|{"SN":"Acc","ST":1574360297,"SV":"0.0"}|
|00001 |1574360355|{"SN":"Acc","ST":1574360297,"SV":"0.0"}|
|00001 |1574360355|{"SN":"Acc","ST":1574360297,"SV":"0.0"}|
|00001 |1574360355|{"SN":"Acc","ST":1574360298,"SV":"0.0"}|
+------+----------+---------------------------------------+

我尝试了以下,但有些行超过了 2MB 的大小,无法写入 cosmosDB:

val newDF = df.groupBy($"ID", $"TIME").agg(collect_list($"SGNL").as("SGNL"))

输出如下,将 n 行合并为 SGNL 列的一行

+------+----------+---------------------------------------------------------------------------------------------------------------------------------------------------------------+
|ID    |TIME      |SGNL                                                                                                                                                           |
+------+----------+---------------------------------------------------------------------------------------------------------------------------------------------------------------+
|00001 |1574360355|{"SN":"Acc","ST":1574360296,"SV":"0.0"},{"SN":"Acc","ST":1574360296,"SV":"0.0"},{"SN":"Acc","ST":1574360296,"SV":"0.0"},.......................................|
|00002 |1574360355|{"SN":"Acc","ST":1574360297,"SV":"0.0"},{"SN":"Acc","ST":1574360297,"SV":"0.0"},{"SN":"Acc","ST":1574360298,"SV":"0.0"},.......................................| 
|00003 |1574360355|{"SN":"Acc","ST":1574360297,"SV":"0.0"},{"SN":"Acc","ST":1574360297,"SV":"0.0"},{"SN":"Acc","ST":1574360298,"SV":"0.0"},.......................................| 
|00004 |1574360355|{"SN":"Acc","ST":1574360297,"SV":"0.0"},{"SN":"Acc","ST":1574360297,"SV":"0.0"},{"SN":"Acc","ST":1574360298,"SV":"0.0"},.......................................|                                        |
+------+----------+---------------------------------------------------------------------------------------------------------------------------------------------------------------+

格式如下:

+------+----------+--------------------------------------------+
|ID    |TIME      |SGNL                                        |
+------+----------+--------------------------------------------+
|00001 |1574360355|{1st ROW},{2nd ROW},......{500th ROW}       |
|00001 |1574360355|{501st ROW},{502nd ROW},......{1000th ROW}  |
|00001 |1574360355|{1001st ROW},{1002nd ROW},......{1500th ROW}|
|00001 |1574360355|{1501st ROW},{1502nd ROW},......{2000th ROW}|
|..............................................................|
|..............................................................|

|00002 |1574360355|{1st ROW},{2nd ROW},......{500th ROW}       |
|00002 |1574360355|{501st ROW},{502nd ROW},......{1000th ROW}  |
|00002 |1574360355|{1001st ROW},{1002nd ROW},......{1500th ROW}|
|00002 |1574360355|{1501st ROW},{1502nd ROW},......{2000th ROW}|
|..............................................................|
+------+----------+---------------------------------------------+

我试图实现这一点,但我只能收集前 n 行,其余行将被忽略。 它不应该只收集前 500 行,而是以 500 的间隔收集所有行。 有什么想法可以实现吗?

【问题讨论】:

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


    【解决方案1】:

    您可以使用monitonicallyIncreasingId 构建索引,将其除以 500 并按该索引分组。它将前 500 行放在一起,接下来的 500 行放在一起,依此类推。

    df.withColumn("fancy_id", floor(monotonicallyIncreasingId / 500))
      .groupBy("fancy_id")
      .agg(collect_list($"SGNL").as("SGNL"))
      .drop("fancy_id") // if you want to get rid of the artificial id.
    

    如果您不想混合 ID 列,可以使用 groupBy("ID", "fancy_id")。

    然而,每个 ID 的第一组的大小不一定是 500。例如,您最终会得到类似:(id1, 500 elements), (id1, 320 elements), (id2, 180 elements), (id2, 500 elements), (id2, 500 element), (id2, 50elements), (id3, 450 elements)...

    如果您更喜欢 (id1, 500 elements), (id1, 320 elements), (id2, 500 elements), (id2, 500 elements), (id2, 500 element), (id2, 10 elements), (id3, 500 elements), (id3, 5 elements)... 之类的东西,其中第一组总是有 500 个元素,您可以使用窗口:

    val w = Window.partitionBy('ID).orderBy('fancy_id)
    df.withColumn("fancy_id", monotonicallyIncreasingId)
      .withColumn("rank", rank() over w)
      .groupBy($"ID", floor($"rank" / 500))
      .agg(collect_list($"SGNL").as("SGNL"))
    

    【讨论】:

    • 这是否会根据任何分区 ID 将 500 行放在一起,还是只会先放置 500 行,然后再放置 500 行,依此类推...?
    • 如果只使用.groupBy("fancy_id"),则分组将基于当前数据的顺序。不同的连续 ID 最终可能会在一起。如果您使用.groupBy("ID", "fancy_id"),不同的ID 将无法组合在一起。然而,组的大小上限为 500,但第一组的大小不一定是 500。如果您需要第一组的大小为 500,则可以使用窗口。如果这是你需要的,我可以更新我的答案。
    • 我目前面临上述实现的问题。我将间隔设置为 5000。对于 5001 条记录,它们被分组为 4999 和 2。对于 10001 条记录,它们被分组为 4999,5000 和 2。对此有什么建议吗?
    • 例如,您可以使用monotonicallyIncreasingId mod 5000 代替除法。它将在 5000 个分区之间平均共享记录。使用除法的目的是硬性限制分区的最大大小(不限制其数量)。
    • 你到底想要什么?
    猜你喜欢
    • 1970-01-01
    • 2019-07-14
    • 2014-03-08
    • 1970-01-01
    • 1970-01-01
    • 2020-09-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多