【发布时间】:2020-10-01 18:08:51
【问题描述】:
我有一个示例 DF,其中包含这样的重复行:
+-------------------+--------------------+----+-----------+-------+----------+
|ID |CL_ID |NBR |DT |TYP |KEY |
+--------------------+--------------------+----+-----------+-------+----------+
|1000031075_20190422 |10017157594301072477|10 |2019-04-24 |N |0000000000|
|1000031075_20190422 |10017157594301072477|10 |2019-04-24 |N |0000000000|
|1006473016_20190421 |10577157412800147475|11 |2019-04-21 |N |0000000000|
|1006473016_20190421 |10577157412800147475|11 |2019-04-21 |N |0000000000|
+--------------------+--------------------+----+-----------+-------+----------+
val w = Window.partitionBy($"ENCOUNTER_ID")
使用上面的 Spark Window 分区,是否可以选择不同的行?我期望输出 DF 为:
+-------------------+--------------------+----+-----------+-------+----------+
|ID |CL_ID |NBR |DT |TYP |KEY |
+--------------------+--------------------+----+-----------+-------+----------+
|1000031075_20190422 |10017157594301072477|10 |2019-04-24 |N |0000000000|
|1006473016_20190421 |10577157412800147475|11 |2019-04-21 |N |0000000000|
+--------------------+--------------------+----+-----------+-------+----------+
我不想使用DF.DISTINCT 或DF.DROPDUPLICATES,因为它会涉及洗牌。
我不喜欢使用滞后或超前,因为实时无法保证行的顺序。
【问题讨论】:
-
窗口函数也需要改组...所以我会使用
dropDuplicates("ID"),distinct更昂贵,因为它比较所有列
标签: apache-spark apache-spark-sql