如何在PySpark中按多条件过滤行:保留重复Product_Number的New记录
解决PySpark DataFrame重复Product_Number的过滤需求
需求回顾
对重复出现的Product_Number仅保留Condition为New的记录,非重复的记录保持原样。
问题分析
你已经正确找出了重复的Product_Number列表,但过滤条件存在两个核心问题:
- 错误将数据中字符串类型的
Null当成了SQL的null值,使用了Condition.isNull()判断,实际应使用Condition == 'Null' - 过滤逻辑不符合需求,没有同时覆盖「非重复记录全保留」和「重复记录仅留New」两个规则
方法一:基于已获取的重复列表过滤(修正版)
先确认重复列表的获取逻辑是正确的:
numbers = df.select('Product_Number')\ .groupBy('Product_Number')\ .count()\ .where('count > 1')\ .select('Product_Number')\ .rdd.flatMap(lambda x: x)\ .collect()
然后使用正确的过滤条件:要么Product_Number不在重复列表中(全保留),要么在重复列表中且Condition为New
filtered_df = df.filter( (~df.Product_Number.isin(numbers)) | (df.Product_Number.isin(numbers) & (df.Condition == 'New')) )
你之前过滤失败的原因
- 第一个filter:
~(isin(numbers)) & ~((Condition == 'Old') | (Condition.isNull())),会过滤掉非重复列表中Condition为Old或Null的记录(比如Product_Number=2和3),不符合需求 - 第二个filter:
~((isin(numbers)) & ((Condition == 'Old') | (Condition.isNull()))),逻辑上等价于「非重复记录全保留,或重复记录中Condition不是Old/Null」,但因为错误用了Condition.isNull()(实际是字符串Null),导致判断失效 - 第三个filter:
~isin(numbers) & (Condition == 'New'),会过滤掉非重复列表中Condition不是New的记录,不符合需求
方法二:使用窗口函数(更高效,无需Collect到Driver)
当数据量较大时,将重复列表Collect到Driver端会有性能问题,推荐用窗口函数实现:
- 按
Product_Number分组,给每条记录排序:New优先级最高,Old次之,Null最低 - 取每个分组的第一条记录
from pyspark.sql import Window from pyspark.sql.functions import row_number, when # 定义窗口规则:分组后按Condition优先级排序 window_spec = Window.partitionBy('Product_Number').orderBy( when(df.Condition == 'New', 1) .when(df.Condition == 'Old', 2) .when(df.Condition == 'Null', 3) .asc() ) # 添加行号并过滤出每组第一条记录 filtered_df = df.withColumn('row_num', row_number().over(window_spec))\ .filter(df.row_num == 1)\ .drop('row_num')
这种方法无需将数据拉到Driver端,适合大数据场景,逻辑也更直观。
内容的提问来源于stack exchange,提问作者enas dyo
相关产品推荐
相关产品推荐

