PySpark按条件删除相似行:移除同组同日期下含1的0值行
解决PySpark DataFrame的分组过滤需求
刚好最近处理过类似的PySpark数据清洗需求,针对你提出的场景——按id和date分组,同一分组内同时存在value=1和value=0时只留value=1的行;如果分组里只有value=0就保留该行,用窗口函数就能高效实现,而且逻辑清晰,不会产生多余的shuffle开销。
先还原测试数据
首先先把你给出的示例DataFrame搭起来:
from pyspark.sql import SparkSession from pyspark.sql.window import Window import pyspark.sql.functions as F spark = SparkSession.builder.appName("FilterValueDemo").getOrCreate() df = spark.createDataFrame( [("A1", "2016-10-01", 1), ("A1", "2016-10-01", 0), ("A1", "2016-10-05", 1), ("A3", "2016-10-06", 1), ("A3", "2016-10-07", 0)], ["id", "date", "value"] )
核心逻辑:用窗口函数标记分组特征
我们需要给每个(id, date)分组打个标记:这个组里有没有value=1的行?这里用Window.partitionBy("id", "date")定义分组窗口,然后用F.max("value")来判断——因为只要组里有1,max结果就是1;全是0的话max就是0:
# 定义窗口:按id和date分组 window_spec = Window.partitionBy("id", "date") # 添加标记列:该分组是否存在value=1 df_with_flag = df.withColumn( "has_value_1", F.max("value").over(window_spec) )
这一步之后,每行都会带上对应分组的has_value_1标记,方便后续过滤。
过滤得到目标结果
现在就可以根据标记和原value来筛选行:
- 如果
has_value_1=1,只留value=1的行(把同组的0删掉) - 如果
has_value_1=0,直接保留value=0的行(因为组里只有它)
对应的过滤代码如下:
result_df = df_with_flag.filter( # 两种符合条件的情况用逻辑或连接 (F.col("has_value_1") == 1) & (F.col("value") == 1) | (F.col("has_value_1") == 0) & (F.col("value") == 0) ).drop("has_value_1") # 删掉临时标记列 # 查看结果 result_df.show()
最终输出
运行后就会得到你想要的结果:
+---+----------+-----+ | id| date|value| +---+----------+-----+ | A1|2016-10-01| 1| | A1|2016-10-05| 1| | A3|2016-10-06| 1| | A3|2016-10-07| 0| +---+----------+-----+
为什么用窗口函数?
相比先分组聚合再join回原表的方式,窗口函数只需要一次shuffle操作,性能更优;而且不用改变原数据的行结构,逻辑也更直观,不容易出错。
内容的提问来源于stack exchange,提问作者Starbucks
相关产品推荐
相关产品推荐

