Spark Dataframe如何按自定义Rank区间分组聚合求和
Spark Dataframe 自定义区间分桶聚合实现方案
核心思路
直接通过Spark原生的when+otherwise条件函数生成自定义区间标签,再按标签分组聚合即可。该方案无UDF额外开销,适配所有不规则区间划分需求,是性能最优的实现方式。
前置依赖导入
Scala 版本
import org.apache.spark.sql.functions.{sum, when, col, concat, lit, floor}
PySpark 版本
from pyspark.sql.functions import sum, when, col, concat, lit, floor
代码实现(按1单独分组、2-5分组、100+分组规则示例)
Scala 版本
val resultDf = df // 生成自定义区间标签 .withColumn("rank_range", when(col("Rank") === 1, "1") .when(col("Rank").between(2, 5), "2-5") // 可按需新增其他区间规则,比如6-99分组可加下一行 // .when(col("Rank").between(6, 99), "6-99") .when(col("Rank") >= 100, "100+") .otherwise("其他") // 兜底规则避免空值 ) // 按区间分组聚合 .groupBy("rank_range") .agg(sum("Col1").alias("sum_col1")) // 按区间逻辑顺序排序输出 .orderBy( when(col("rank_range") === "1", 1) .when(col("rank_range") === "2-5", 2) .when(col("rank_range") === "6-99", 3) .when(col("rank_range") === "100+", 4) ) // 输出结果 resultDf.show()
PySpark 版本
result_df = df\ .withColumn("rank_range", when(col("Rank") == 1, "1") .when(col("Rank").between(2, 5), "2-5") # 可按需新增其他区间规则 # .when(col("Rank").between(6, 99), "6-99") .when(col("Rank") >= 100, "100+") .otherwise("其他") )\ .groupBy("rank_range")\ .agg(sum("Col1").alias("sum_col1"))\ .orderBy( when(col("rank_range") == "1", 1) .when(col("rank_range") == "2-5", 2) .when(col("rank_range") == "6-99", 3) .when(col("rank_range") == "100+", 4) ) result_df.show()
其他场景适配
- 如果你已经提前计算好了单Rank聚合的结果,只需要把上述代码中的
df替换为你的单Rank聚合结果表,把聚合逻辑改成sum("sum(Col1)")(对应你原来聚合结果的求和列名)即可。 - 如果是固定步长的规则区间,比如每5个Rank为一组,可以直接用数值计算生成标签,不需要逐个写
when条件,示例:// 步长5的规则区间生成逻辑 .withColumn("rank_range", concat(floor(col("Rank")/5)*5, lit("-"), floor(col("Rank")/5)*5 + 4) )
内容的提问来源于stack exchange,提问作者Sash9
相关产品推荐
相关产品推荐

