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
相关产品推荐
相关产品推荐

