Spark SQL大数据集下dense_rank计算错误问题排查
我正在对一张约100万行的表计算dense_rank,针对TotalEmployees的不同区间分别计算百分位,输出结果共100万条,但排名不符合预期——数据集中的最大值未被分配最高排名或100分位。已知Spark SQL通过数据分区/块加速处理,请问这是否是问题根源?
排名函数:dense_rank
| pageViewsCount | Expected Rank | Wrong Rank |
|---|---|---|
| 99 | 6 | 2 |
| 44 | 5 | 1 |
| 31 | 4 | 5 |
| 25 | 3 | 4 |
| 21 | 2 | 3 |
| 21 | 2 | 2 |
| 10 | 1 | 1 |
用户使用的SQL代码:
select named_struct('JobName','Traffic') AS headers, named_struct('data',named_struct('companyId',CompanyId,'trafficPercentile',percentile,'pageViewsCount',pageViewsCount,'category','Traffic')) as payload from ( select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 1 and TotalEmployees <= 10) as dist_cnt from pv where TotalEmployees >= 1 and TotalEmployees <= 10 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 11 and TotalEmployees <= 50) as dist_cnt from pv where TotalEmployees >= 11 and TotalEmployees <= 50 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 51 and TotalEmployees <= 200) as dist_cnt from pv where TotalEmployees >= 51 and TotalEmployees <= 200 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 201 and TotalEmployees <= 500) as dist_cnt from pv where TotalEmployees >= 201 and TotalEmployees <= 500 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 501 and TotalEmployees <= 1000) as dist_cnt from pv where TotalEmployees >= 501 and TotalEmployees <= 1000 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 1001 and TotalEmployees <= 5000) as dist_cnt from pv where TotalEmployees >= 1001 and TotalEmployees <= 5000 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 5001 and TotalEmployees <= 10000) as dist_cnt from pv where TotalEmployees >= 5001 and TotalEmployees <= 10000 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 10001 and TotalEmployees <= 50000) as dist_cnt from pv where TotalEmployees >= 10001 and TotalEmployees <= 50000 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 50001 and TotalEmployees <= 100000) as dist_cnt from pv where TotalEmployees >= 50001 and TotalEmployees <= 100000 ) union select *,(dense_rank/dist_cnt) * 100 as percentile from( select CompanyId,pageViewsCount,(dense_rank() OVER (ORDER BY pageViewsCount)) as dense_rank,TotalEmployees,(select count(distinct(pageViewsCount)) from pv where TotalEmployees >= 100001) as dist_cnt from pv where TotalEmployees >= 100001 ) order by TotalEmployees )
问题根源:不是Spark分区,是SQL写法和数据类型的问题
1. 数据类型或排序方向错误导致排名混乱
从错误结果来看,pageViewsCount=44排名为1、99排名为2,这完全符合字符串字典序的排序逻辑("44"在"99"之前)。如果pageViewsCount是字符串类型而非数值类型,Spark会按字典序排序而非数值大小,直接导致排名错误。
另外,若你希望最大值对应100分位,即使数据类型正确,当前的升序排序(ORDER BY pageViewsCount)会让最大值的dense_rank等于该区间的不同值数量dist_cnt,此时(dense_rank/dist_cnt)*100确实能得到100分位,但如果你的预期是最高排名的数值为1(比如让最大值的rank为1),则需要改为降序排序:dense_rank() OVER (ORDER BY pageViewsCount DESC)。
2. 冗余的UNION写法易出错,应改用PARTITION BY
当前将每个TotalEmployees区间拆分为独立子查询再UNION的写法,不仅冗余,还容易因区间边界错误导致数据被分到错误分组,进而出现排名混乱。正确的做法是用CASE WHEN将TotalEmployees映射为分组,再通过PARTITION BY按分组计算排名和分位数,代码更简洁高效:
select named_struct('JobName','Traffic') AS headers, named_struct('data',named_struct('companyId',CompanyId,'trafficPercentile',percentile,'pageViewsCount',pageViewsCount,'category','Traffic')) as payload from ( select CompanyId, pageViewsCount, TotalEmployees, -- 按TotalEmployees划分区间分组 CASE WHEN TotalEmployees BETWEEN 1 AND 10 THEN '1-10' WHEN TotalEmployees BETWEEN 11 AND 50 THEN '11-50' WHEN TotalEmployees BETWEEN 51 AND 200 THEN '51-200' WHEN TotalEmployees BETWEEN 201 AND 500 THEN '201-500' WHEN TotalEmployees BETWEEN 501 AND 1000 THEN '501-1000' WHEN TotalEmployees BETWEEN 1001 AND 5000 THEN '1001-5000' WHEN TotalEmployees BETWEEN 5001 AND 10000 THEN '5001-10000' WHEN TotalEmployees BETWEEN 10001 AND 50000 THEN '10001-50000' WHEN TotalEmployees BETWEEN 50001 AND 100000 THEN '50001-100000' WHEN TotalEmployees >= 100001 THEN '100001+' END as employee_group, -- 按分组降序计算dense_rank,最大值rank为1 dense_rank() OVER (PARTITION BY employee_group ORDER BY pageViewsCount DESC) as dense_rank, -- 按分组计算不同pageViewsCount的数量 COUNT(DISTINCT pageViewsCount) OVER (PARTITION BY employee_group) as dist_cnt, -- 计算百分位 (dense_rank / COUNT(DISTINCT pageViewsCount) OVER (PARTITION BY employee_group)) * 100 as percentile from pv where TotalEmployees >= 1 ) order by TotalEmployees
3. 额外验证点
- 确认
pageViewsCount的数据类型为数值类型(int/bigint/double),若为字符串需先转换:CAST(pageViewsCount AS INT) - 检查每个TotalEmployees区间的边界是否正确,避免数据漏分或错分
内容的提问来源于stack exchange,提问作者harshit rathod

