PySpark中按ID筛选列匹配行 无匹配则保留分组首行
实现方案
你原来的写法只过滤了全局满足numb == real的行,没有处理分组内无匹配行的场景,可通过窗口函数加标记的方式实现需求:
PySpark 实现代码
from pyspark.sql.functions import col, when, max, row_number from pyspark.sql.window import Window # 定义分组窗口,按ID分区,按numb升序排序(如需按原始数据顺序,可提前新增monotonically_increasing_id作为排序键) window_spec = Window.partitionBy("ID").orderBy("numb") # 新增两个标记列:标记分组是否有匹配行、分组内行号 df = df.withColumn("has_match", max(when(col("numb") == col("real"), 1).otherwise(0)).over(window_spec)) \ .withColumn("row_num", row_number().over(window_spec)) # 按规则筛选 df_result = df.filter( # 有匹配的分组,保留所有匹配行 (col("has_match") == 1) & (col("numb") == col("real")) | # 无匹配的分组,保留首行 (col("has_match") == 0) & (col("row_num") == 1) ).drop("has_match", "row_num") # 输出结果验证 df_result.show()
运行后输出结果和预期一致:
+------+----+----+ | ID|numb|real| +------+----+----+ |213412|2008|2008| |213410|2009|null| |393859|2017|2017| |393859|2021|2021| +------+----+----+
Pandas 实现代码
如果你用的是Pandas,可以用groupby apply的方式实现:
import pandas as pd import numpy as np def filter_group(g): # 先找当前分组匹配的行 match_rows = g[g["numb"] == g["real"]] if not match_rows.empty: return match_rows # 没有匹配返回首行 return g.head(1) df_result = df.groupby("ID", group_keys=False).apply(filter_group).reset_index(drop=True)
内容的提问来源于stack exchange,提问作者BADS
相关产品推荐
相关产品推荐

