【发布时间】: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