【问题标题】:Spark Dataframe - How to keep only latest record for each group based on ID and Date? [duplicate]Spark Dataframe - 如何根据 ID 和日期仅保留每个组的最新记录? [复制]
【发布时间】:2020-05-10 04:08:57
【问题描述】:

我有一个数据框:

DF:

1,2016-10-12 18:24:25
1,2016-11-18 14:47:05
2,2016-10-12 21:24:25
2,2016-10-12 20:24:25
2,2016-10-12 22:24:25
3,2016-10-12 17:24:25

如何只保留每个组的最新记录? (上面有3组(1,2,3))。

结果应该是:

1,2016-11-18 14:47:05
2,2016-10-12 22:24:25
3,2016-10-12 17:24:25

还努力提高效率(例如,在具有 1 亿条记录的中等集群上在短短几分钟内完成),因此应以最有效和正确的方式进行排序/排序(如果需要)。.

【问题讨论】:

标签: dataframe date apache-spark pyspark


【解决方案1】:

您可以使用here 中描述的窗口函数来处理以下情况:

scala> val in = Seq((1,"2016-10-12 18:24:25"),
     | (1,"2016-11-18 14:47:05"),
     | (2,"2016-10-12 21:24:25"),
     | (2,"2016-10-12 20:24:25"),
     | (2,"2016-10-12 22:24:25"),
     | (3,"2016-10-12 17:24:25")).toDF("id", "ts")
in: org.apache.spark.sql.DataFrame = [id: int, ts: string]
scala> import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.expressions.Window

scala> val win = Window.partitionBy("id").orderBy("ts desc")
win: org.apache.spark.sql.expressions.WindowSpec = org.apache.spark.sql.expressions.WindowSpec@59fa04f7

scala> in.withColumn("rank", row_number().over(win)).where("rank == 1").show(false)
+---+-------------------+----+
| id|                 ts|rank|
+---+-------------------+----+
|  1|2016-11-18 14:47:05|   1|
|  3|2016-10-12 17:24:25|   1|
|  2|2016-10-12 22:24:25|   1|
+---+-------------------+----+

【讨论】:

  • 你如何在 Pyspark 中编写这个?
  • 我包含的链接包含 Pyspark 中的几个示例
  • 谢谢,除了最后一行代码之外,我可以转移所有代码:in.withColumn("rank", row_number().over(win)).where('rank === 1).show(false)
【解决方案2】:

你必须使用窗口功能。

http://spark.apache.org/docs/latest/api/python/pyspark.sql.html?highlight=window#pyspark.sql.Window

您必须按组和 OrderBy 时间对窗口进行分区,下面的 pyspark 脚本完成工作

from pyspark.sql.functions import *
from pyspark.sql.window import Window

schema = "Group int,time timestamp "

df = spark.read.format('csv').schema(schema).options(header=False).load('/FileStore/tables/Group_window.txt')


w = Window.partitionBy('Group').orderBy(desc('time'))
df = df.withColumn('Rank',dense_rank().over(w))

df.filter(df.Rank == 1).drop(df.Rank).show()


+-----+-------------------+
|Group|               time|
+-----+-------------------+
|    1|2016-11-18 14:47:05|
|    3|2016-10-12 17:24:25|
|    2|2016-10-12 22:24:25|
+-----+-------------------+ ```





【讨论】:

  • 你应该使用 row_number() 而不是 dense_rank() 因为dense_rank() 为“并列”行提供相同的排名
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-01-23
  • 2015-11-07
  • 1970-01-01
  • 2020-08-03
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多