You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Spark SQL中对相似数值进行分组?

Spark SQL 数值相似分组及平均值计算方案

核心思路是通过LAG()窗口函数获取前一行的数值,判断是否触发分组切换,再累加切换标志生成组索引,最后分组求平均。以下是具体步骤和代码:

步骤分解

  1. 获取前一行数值:用LAG()按idx排序,拿到当前行的前一行col1值。
  2. 生成分组切换标志:对比当前值与前一行值的相对差异,超过阈值(比如5%)则标记为1(切换分组),否则为0。同时处理第一行、前一行值为0的特殊情况。
  3. 生成组索引:累加切换标志,得到递增的组ID(即示例中的col2)。
  4. 计算组内平均值:按组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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 12:40:26