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的行)
实现思路
核心分两步:
- 先预处理所有
type='gift'的行,按id,nr,prod分组,提取每个分组里pg_nr最大的那一行的var_1和var_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
相关产品推荐
相关产品推荐

