如何在Spark SQL中对相似数值进行分组?
Spark SQL 数值相似分组及平均值计算方案
核心思路是通过LAG()窗口函数获取前一行的数值,判断是否触发分组切换,再累加切换标志生成组索引,最后分组求平均。以下是具体步骤和代码:
步骤分解
- 获取前一行数值:用
LAG()按idx排序,拿到当前行的前一行col1值。 - 生成分组切换标志:对比当前值与前一行值的相对差异,超过阈值(比如5%)则标记为1(切换分组),否则为0。同时处理第一行、前一行值为0的特殊情况。
- 生成组索引:累加切换标志,得到递增的组ID(即示例中的
col2)。 - 计算组内平均值:按组ID分组,计算
col1的平均值。
完整SQL代码
假设你的视图名为my_view,阈值设为5%,代码如下:
-- 步骤1:获取前一行的col1值 WITH step1 AS ( SELECT idx, col1, LAG(col1) OVER(ORDER BY idx) AS prev_col1 FROM my_view ), -- 步骤2:生成分组切换标志 step2 AS ( SELECT idx, col1, prev_col1, CASE WHEN prev_col1 IS NULL THEN 0 -- 第一行无前置数据,不切换分组 WHEN prev_col1 = 0 AND col1 != 0 THEN 1 -- 前一行是0,当前行非0,切换分组 WHEN ABS((col1 - prev_col1)/prev_col1) > 0.05 THEN 1 -- 差异超过5%,切换分组 ELSE 0 -- 差异在阈值内,保留在当前组 END AS group_flag FROM step1 ), -- 步骤3:累加标志生成组索引 step3 AS ( SELECT idx, col1, SUM(group_flag) OVER(ORDER BY idx) AS col2 FROM step2 ) -- 步骤4:计算每组的col1平均值 SELECT col2 AS group_id, ROUND(AVG(col1), 4) AS avg_col1 -- 保留4位小数,可按需调整 FROM step3 GROUP BY col2 ORDER BY group_id;
代码说明
- 处理特殊情况:针对前一行值为0的场景单独判断,避免出现除以0的运算错误。
- 相对差异计算:用
ABS((col1 - prev_col1)/prev_col1)计算相对变化比例,确保是基于前一行数值的百分比差异。 - 组索引生成:通过
SUM(group_flag) OVER(ORDER BY idx)累加切换标志,每次切换时组ID自动加1,符合你需要的递增效果。
示例结果
按你的测试数据运行后,step3会生成你需要的col2列:
idx | col1 | col2 ------------------- 0 | 0.0 | 0 1 | 0.0 | 0 2 | 0.384 | 1 3 | 0.383 | 1 4 | 0.382 | 1 5 | 0.384 | 1 6 | 0.518 | 2 7 | 0.524 | 2 8 | 0.520 | 2 9 | 0.250 | 3
最终分组平均值结果:
group_id | avg_col1 ------------------- 0 | 0.0000 1 | 0.3832 2 | 0.5207 3 | 0.2500
内容的提问来源于stack exchange,提问作者conni
相关产品推荐
相关产品推荐

