Spark SQL/PySpark替代SQL Server关联子查询实现NgId分组填充
基于NewGroupId分组填充空NgId的Spark实现
需求
按NewGroupId字段分组,将分组内NgId为空的记录,用同组内非空的NgId值覆盖。
输入示例
Name GroupId Processed NewGroupId NgId Mike 1 N 9 NULL Mikes 1 N 9 NULL Miken 5 Y 9 5 Mikel 5 Y 9 5
输出示例
Name GroupId Processed NewGroupId NgId Mike 1 N 9 5 Mikes 1 N 9 5 Miken 5 Y 9 5 Mikel 5 Y 9 5
原SQL Server的关联子查询写法无法在Spark SQL中运行:
SELECT Name,groupid,IsProcessed,ngid, CASE WHEN ngid IS NULL THEN COALESCE((SELECT top 1 ngid FROM temp D WHERE D.NewGroupId = T.NewGroupId AND D.ngid IS NOT NULL ), null) ELSE ngid END AS ngid FROM temp T
替代实现方案
1. Spark SQL 写法
利用窗口函数可简洁实现,因同分组内非空NgId值一致,推荐两种写法:
用MAX窗口函数(更简洁)
SELECT Name, GroupId, Processed, NewGroupId, COALESCE(NgId, MAX(NgId) OVER (PARTITION BY NewGroupId)) AS NgId FROM temp;
用FIRST_VALUE窗口函数(忽略空值)
SELECT Name, GroupId, Processed, NewGroupId, COALESCE( NgId, FIRST_VALUE(NgId IGNORE NULLS) OVER (PARTITION BY NewGroupId) ) AS NgId FROM temp;
2. PySpark DataFrame 写法
提供两种常用实现方式:
方式一:窗口函数直接填充
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, coalesce, max, first spark = SparkSession.builder.appName("FillNgIdNulls").getOrCreate() # 假设df为输入数据的DataFrame window_spec = Window.partitionBy("NewGroupId") # 用MAX函数填充 df_filled = df.withColumn( "NgId", coalesce(col("NgId"), max(col("NgId")).over(window_spec)) ) # 或者用FIRST_VALUE忽略空值填充 df_filled = df.withColumn( "NgId", coalesce(col("NgId"), first(col("NgId"), ignorenulls=True).over(window_spec)) ) df_filled.show()
方式二:分组聚合后关联
先提取每个分组的非空NgId,再关联原表填充:
from pyspark.sql.functions import first # 获取每个NewGroupId对应的非空NgId值 grouped_ref = df.filter(col("NgId").isNotNull()) \ .groupBy("NewGroupId") \ .agg(first("NgId").alias("FilledNgId")) # 关联原表并填充空值 df_filled = df.join(grouped_ref, on="NewGroupId", how="left") \ .withColumn("NgId", coalesce(col("NgId"), col("FilledNgId"))) \ .drop("FilledNgId") df_filled.show()
内容的提问来源于stack exchange,提问作者Adi
相关产品推荐
相关产品推荐

