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

SparkSQL中基于分组匹配行更新指定字段的实现方案

SparkSQL中基于分组匹配行更新指定字段的实现方案

嘿,我来帮你搞定这个SparkSQL的更新需求!先咱们把规则再捋一遍,确保没理解错:

  • 按id和nr分组
  • 找到组内满足(type = 'sold' OR page = 'bag')的行,当这些行的prod和同组里type = 'gift'的行的prod一致时
  • 用这些gift行里最大pg_nr对应的var_1和var_2,更新目标行的这两个字段(而且看你的输出,应该是只更新原本var_1/var_2为null的行)

实现思路

核心分两步:

  1. 先预处理所有type='gift'的行,按id,nr,prod分组,提取每个分组里pg_nr最大的那一行的var_1和var_2——这就是我们要用来更新的参考值
  2. 将原表和预处理后的参考数据关联,对符合条件的行执行更新

方案1:使用MERGE INTO(Spark 2.4+支持)

MERGE INTO是SparkSQL里用来做UPSERT的官方语法,适合直接更新原表:

-- 第一步:预处理gift类型行,获取每个(id,nr,prod)组内最大pg_nr对应的var值
WITH gift_prod_vars AS (
    SELECT 
        id,
        nr,
        prod,
        var_1,
        var_2
    FROM (
        SELECT 
            id,
            nr,
            prod,
            var_1,
            var_2,
            -- 按pg_nr降序排,取第一行就是最大pg_nr的记录
            ROW_NUMBER() OVER (PARTITION BY id, nr, prod ORDER BY pg_nr DESC) AS rn
        FROM My_Table
        WHERE type = 'gift'
    ) t
    WHERE rn = 1
)
-- 第二步:关联原表执行更新
MERGE INTO My_Table target
USING (
    SELECT 
        src.id,
        src.nr,
        src.pg_nr,
        gp.var_1 AS new_var_1,
        gp.var_2 AS new_var_2
    FROM My_Table src
    LEFT JOIN gift_prod_vars gp
        ON src.id = gp.id
        AND src.nr = gp.nr
        AND src.prod = gp.prod
    -- 筛选出需要更新的目标行:满足type/sold条件,且原var值为null
    WHERE (src.type = 'sold' OR src.page = 'bag') AND src.var_1 IS NULL
) source
ON target.id = source.id 
   AND target.nr = source.nr 
   AND target.pg_nr = source.pg_nr
WHEN MATCHED THEN
    UPDATE SET var_1 = source.new_var_1, var_2 = source.new_var_2;

方案2:创建新表(兼容低版本Spark)

如果你的Spark版本不支持MERGE INTO,可以用创建新表的方式实现,本质是通过CASE WHEN判断是否需要更新:

WITH gift_prod_vars AS (
    SELECT 
        id,
        nr,
        prod,
        var_1,
        var_2
    FROM (
        SELECT 
            id,
            nr,
            prod,
            var_1,
            var_2,
            ROW_NUMBER() OVER (PARTITION BY id, nr, prod ORDER BY pg_nr DESC) AS rn
        FROM My_Table
        WHERE type = 'gift'
    ) t
    WHERE rn = 1
),
updated_data AS (
    SELECT 
        t.id,
        t.nr,
        t.pg_nr,
        t.type,
        t.list,
        t.page,
        t.prod,
        -- 只更新符合条件且原var为null的行
        CASE 
            WHEN (t.type = 'sold' OR t.page = 'bag') AND gp.var_1 IS NOT NULL AND t.var_1 IS NULL THEN gp.var_1
            ELSE t.var_1
        END AS var_1,
        CASE 
            WHEN (t.type = 'sold' OR t.page = 'bag') AND gp.var_2 IS NOT NULL AND t.var_2 IS NULL THEN gp.var_2
            ELSE t.var_2
        END AS var_2
    FROM My_Table t
    LEFT JOIN gift_prod_vars gp
        ON t.id = gp.id
        AND t.nr = gp.nr
        AND t.prod = gp.prod
)
-- 输出更新后的结果,也可以写成INSERT OVERWRITE INTO 新表名
SELECT * FROM updated_data;

结果验证

运行上面的代码后,就能得到你给出的输出结果:

  • v1,1,pg_nr=6的sold行(prod=SE92)会匹配到gift组里pg_nr=4的SE92记录,把null的var更新为ball和orange
  • 像v1,1,pg_nr=1这种已有var值的行,即使prod匹配,也不会被修改
  • 其他不满足条件的行保持原样

备注:内容来源于stack exchange,提问作者D.Aquar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 10:25:33