Spark SQL分组后累计求和异常:相同count值导致结果不符合预期
解决Spark SQL累计求和(running_sum)在count值相同时合并的问题
原始数据
使用以下代码创建Spark DataFrame并注册为临时视图:
data = [("1","1"), ("1","1"), ("1","1"), ("2","1"), ("2","1"), ("3","1"), ("3","1"), ("4","1"),] df = spark.createDataFrame(data=data, schema=["id","imp"]) df.createOrReplaceTempView("df")
DataFrame内容如下:
+---+---+ | id|imp| +---+---+ | 1| 1| | 1| 1| | 1| 1| | 2| 1| | 2| 1| | 3| 1| | 3| 1| | 4| 1| +---+---+
需求
按id分组,统计每个id的出现次数(count)、累计求和(running_sum)以及所有id的总次数(total_sum)。
当前实现及问题
使用以下SQL查询:
select id, count(id) as count, sum(count(id)) over (order by count(id) desc) as running_sum, sum(count(id)) over () as total_sum from df group by id order by count desc
执行结果:
+---+-----+-----------+---------+ | id|count|running_sum|total_sum| +---+-----+-----------+---------+ | 1| 3| 3| 8| | 2| 2| 7| 8| | 3| 2| 7| 8| | 4| 1| 8| 8| +---+-----+-----------+---------+
问题:当多个id的count值相同时,running_sum会直接合并这些行的count值求和,导致ID2和ID3的running_sum均为7,不符合预期的逐行累计逻辑。
预期结果
+---+-----+-----------+---------+ | id|count|running_sum|total_sum| +---+-----+-----------+---------+ | 1| 3| 3| 8| | 2| 2| 5| 8| | 3| 2| 7| 8| | 4| 1| 8| 8| +---+-----+-----------+---------+
解决方案
问题根源在于窗口函数的排序仅依赖count(id) desc,相同count值的行被视为同一分组,默认的RANGE窗口范围会将这些行全部纳入累计计算。修改方案如下:
- 在
order by中添加id作为次要排序键,确保每行的排序唯一; - 显式指定窗口范围为
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW,明确累计求和从第一行到当前行。
修改后的SQL:
select id, count(id) as count, sum(count(id)) over (order by count(id) desc, id asc rows between unbounded preceding and current row) as running_sum, sum(count(id)) over () as total_sum from df group by id order by count desc, id asc
执行该SQL即可得到符合预期的累计求和结果。
内容的提问来源于stack exchange,提问作者Vishal Balaji
相关产品推荐
相关产品推荐

