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

PySpark按指定列分组选取空字段数量最少的记录实现方法

PySpark 按分组取空值最少记录实现方案

实现逻辑

核心思路是先统计每条记录中除分组字段外的空值数量,再通过窗口函数按分组排序取空值最少的第一条记录,支持自动适配多字段场景,无需手动枚举所有非分组字段。

完整可运行代码

# 导入依赖
from pyspark.sql import Window
import pyspark.sql.functions as F

# 1. 定义分组去重字段
dedup_cols = ['facid','mrn','date','filetimestamp']
# 2. 自动获取需要统计空值的非分组字段(适配20+字段场景,无需手动写全)
count_null_cols = [col for col in spark_df.columns if col not in dedup_cols]
# 3. 新增列统计当前行非分组字段的空值总数
df_with_null_cnt = spark_df.withColumn(
    "null_count",
    sum(F.when(F.col(c).isNull(), 1).otherwise(0) for c in count_null_cols)
)
# 4. 定义窗口:按分组字段分区,空值数量升序排序(空值越少越靠前)
# 若同分组内存在多条空值数量相同的记录,可追加排序字段保证结果稳定,示例追加mdsid排序
window_spec = Window.partitionBy(*dedup_cols).orderBy(F.asc("null_count"), F.asc("mdsid"))
# 5. 窗口内排名,取每组排名第一的记录
df_with_rank = df_with_null_cnt.withColumn("rn", F.row_number().over(window_spec))
final_df = df_with_rank.filter(F.col("rn") == 1).drop("null_count", "rn")

# 输出结果验证
final_df.show()

输出结果验证

运行后输出结果与需求一致:

+-------+----+-------------------+-----------------+------------------+-----+---------+
|  facid| mrn|               date|            mdsid|     filetimestamp|e0300|v0200a14b|
+-------+----+-------------------+-----------------+------------------+-----+---------+
|PMS.RCM|2442|2021-03-15T00:00:00|1a644a8c08800e0b7|202104060445575041|    0|        1|
|PMS.RCM|2442|2021-03-12T00:00:00|f5a79fa8beb64d6e4|202104060445575041| null|     null|
+-------+----+-------------------+-----------------+------------------+-----+---------+

内容的提问来源于stack exchange,提问作者Sisay

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 16:15:03