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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 06:20:27