PySpark DataFrame按ID和Month分组后过滤仅出现一次的行
PySpark 实现双分组过滤单次出现组合的行
你需要保留ID和Month组合出现次数≥2的所有原始行,可通过以下两种方式实现,两种方式都不会修改原始数据的内容,仅做行过滤:
方法1:窗口函数实现(推荐,性能更优)
窗口函数可以在不聚合原始行的前提下,给每个分组标记统计值,刚好适配你的需求:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按ID和Month分组,不需要排序 w = Window.partitionBy("ID", "Month") # 新增分组计数列,过滤计数>1的行后删除冗余的计数列 df_result = df.withColumn("group_count", F.count("*").over(w)) \ .filter(F.col("group_count") > 1) \ .drop("group_count") # 查看结果 df_result.show()
方法2:分组统计后关联实现
如果你更习惯分组统计的逻辑,可以先统计符合要求的组合,再和原表关联过滤:
from pyspark.sql import functions as F # 统计每个ID+Month组合的出现次数,过滤出出现≥2次的有效组合 valid_groups = df.groupBy("ID", "Month").count().filter(F.col("count")>1).select("ID", "Month") # 用内连接保留原表中属于有效组合的所有行 df_result = df.join(valid_groups, on=["ID", "Month"], how="inner") # 查看结果 df_result.show()
之前代码的问题说明
- pandas风格的布尔索引写法不支持PySpark DataFrame,两者的API设计逻辑差异较大,不能直接套用
- 之前写的窗口函数错误:分区字段仅设置了Month,没有同时按ID分区,也没有正确使用count聚合做分组计数,还引用了不存在的变量
Moth
运行上述两种方法得到的结果和你给出的预期输出完全一致。
内容的提问来源于stack exchange,提问作者Meike
相关产品推荐
相关产品推荐

