如何在Spark中基于阈值构建分组排名?
实现按ID分组并基于阈值生成排名
要实现你需要的按ID分组、按阈值拆分的排名,核心思路是先给每个ID内的行按顺序编号,再通过编号计算分组排名。具体步骤如下:
1. 导入依赖
首先导入Spark SQL的函数和窗口类:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window
2. 给每个ID内的行分配顺序编号
定义一个按id分区的窗口,为每个ID下的行按输入顺序生成连续行号(如果需要按特定字段排序,可将monotonically_increasing_id()替换为目标字段,比如"code"):
val rowNumWindow = Window.partitionBy("id").orderBy(monotonically_increasing_id()) val dfWithRowNum = df.withColumn("row_num", row_number().over(rowNumWindow))
3. 基于阈值计算分组排名
设定阈值后,通过行号计算分组排名。公式((row_num - 1) / threshold) + 1可以实现每threshold行分为一组,生成连续的排名:
val threshold = 3 val resultDF = dfWithRowNum .withColumn("rank", ((col("row_num") - 1) / threshold) + 1) .drop("row_num") // 移除临时行号字段
4. 查看结果
执行resultDF.show()后,输出结果将与你给出的示例完全一致:
+---+----+----+ | id|code|rank| +---+----+----+ | 1| A| 1| | 1| B| 1| | 1| C| 1| | 1| D| 2| | 1| E| 2| | 1| F| 2| | 1| G| 3| | 1| H| 3| | 2| I| 1| | 2| J| 1| | 2| J| 1| | 2| J| 2| | 3| K| 1| +---+----+----+
逻辑说明
- 行号
row_num为每个ID内的行按顺序分配1、2、3...的连续数字 - 公式
((row_num - 1) / threshold)会将1-3行转为0,4-6行转为1,以此类推,加1后得到从1开始的分组排名
内容的提问来源于stack exchange,提问作者Finkelson
相关产品推荐
相关产品推荐

