如何将Spark SQL改写为纯Scala(Spark)代码实现每年最受欢迎姓名查询
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._
// 从已注册的临时表读取原始数据 val popularNamesDF = spark.table("popularnames") // 定义窗口规则:按年份分区,同分区内按姓名出现次数倒序排列 val yearWindow = Window.partitionBy("year").orderBy(desc("cnt")) val resultDF = popularNamesDF // 按年份、姓名分组统计出现次数,对应SQL中group by + count逻辑 .groupBy("year", "first_name") .agg(count("*").alias("cnt")) // 给每个年份内的记录打行号,对应SQL中row_number窗口函数逻辑 .withColumn("seqnum", row_number().over(yearWindow)) // 过滤取每个年份排名第一的记录 .filter(col("seqnum") === 1) // 取目标输出字段 .select("year", "first_name")
如果需要保留同一年份出现次数并列第一的所有姓名,将row_number()替换为rank()即可,和原SQL逻辑完全对齐。
内容的提问来源于stack exchange,提问作者olifant71
相关产品推荐
相关产品推荐

