如何使用PySpark按小时获取指定数量的高频LocationID?
获取每小时频率最高的Top2 LocationID
我来帮你搞定这个需求!要从你的Spark DataFrame里提取每小时出现频率最高的前2个locationID,用窗口函数是最直接高效的方案,这也是Spark处理分组TopN问题的标准做法。
完整实现步骤
假设你的DataFrame名称是df,我们可以按下面的步骤操作:
- 首先导入需要的函数模块:
from pyspark.sql import Window from pyspark.sql.functions import row_number, desc
- 定义窗口规则:
我们需要按hour字段分组,然后在每个分组内按frequency降序排序,这样频率最高的location会排在前面:
# 按hour分区,区内按frequency降序排列 window_spec = Window.partitionBy("hour").orderBy(desc("frequency"))
- 添加排名列并过滤Top2:
给每个分组内的行添加排名,然后筛选出排名前2的记录,最后去掉多余的排名列:
top2_locations = df.withColumn("rank", row_number().over(window_spec)) \ .filter("rank <= 2") \ .drop("rank")
- 查看结果:
执行top2_locations.show(),针对你给出的示例数据,输出会是这样:
+----+----------+---------+ |hour|locationID|frequency| +----+----------+---------+ | 0 | 1 | 20 | | 0 | 2 | 11 | | 1 | 3 | 32 | | 1 | 1 | 22 | +----+----------+---------+
额外说明
- 如果遇到同一小时内多个location的
frequency相同的情况,row_number()会给它们分配不同的排名(比如两个location频率都是20,一个排1,一个排2)。如果你希望并列的location都被保留(比如两个频率20的都算Top2),可以把row_number()换成rank()或者dense_rank():rank():相同频率会得到相同排名,后续排名会跳号(比如1,1,3)dense_rank():相同频率得到相同排名,后续排名不会跳号(比如1,1,2)
根据你的需求选择合适的排名函数就好啦!
内容的提问来源于stack exchange,提问作者Daniele Isoni
相关产品推荐
相关产品推荐

