每日更新行的去重问题:如何避免重复统计同一客户行?
问题背景
有一个每日更新的DataFrame,包含Customer ID、status和date字段。部分客户每日会收到状态更新,部分则不会;客户状态可能在多日之间在no和yes之间切换。
现有筛选逻辑
筛选2022年10月状态为yes的行的代码如下:
df = df \ .select('id','status','date') \ .filter( (col('date') >= '2022-10-01') & (col('date') <= '2022-10-31') & (col(status) == "yes"))
需求说明
需要生成第二个筛选结果集,要求完全排除上述查询中出现过的所有ID。例如:若ID“123”在10月有过yes状态的记录,即便该ID存在no状态的行,也不能将其计入“no”部分的统计。
尝试的错误实现及报错
尝试用Window函数按ID分区创建标记以排除目标ID,但执行报错,PySpark提示窗口函数不支持该表达式:
partition = Window.partitionBy("id").orderBy("date") df = df \ .withColumn("results", when((col("status") == "approved").over(partition), '0') .otherwise("1"))
报错信息:
Py4JJavaError: An error occurred while calling o808.withColumn. : org.apache.spark.sql.AnalysisException: Expression '(result_decisaofinal#8593 = APROVA)' not supported within a window function.;;
解决方案
方法1:先获取目标ID集合再过滤
先提取2022年10月状态为yes的所有唯一ID,再通过左反连接排除这些ID(可按需添加日期过滤条件):
from pyspark.sql import functions as F # 获取10月有过yes状态的所有唯一ID yes_ids = df.filter( (F.col('date') >= '2022-10-01') & (F.col('date') <= '2022-10-31') & (F.col('status') == "yes") ).select('id').distinct() # 筛选排除上述ID的结果集 no_result_df = df.join(yes_ids, on='id', how='left_anti')
left_anti连接会保留左表中未出现在右表的所有行,恰好满足“完全排除有过yes记录的ID”的需求。
方法2:用窗口函数标记ID是否存在yes状态
之前的错误在于直接将布尔表达式附加.over(partition),窗口函数需配合聚合函数使用。正确写法如下:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 按ID分区,标记该ID在10月是否有过yes状态 window_spec = Window.partitionBy("id") df_with_flag = df.withColumn( "has_yes_in_oct", F.max( F.when( (F.col('date') >= '2022-10-01') & (F.col('date') <= '2022-10-31') & (F.col('status') == "yes"), 1 ).otherwise(0) ).over(window_spec) ) # 筛选出未出现过yes状态的ID对应的行 no_result_df = df_with_flag.filter(F.col("has_yes_in_oct") == 0)
内容的提问来源于stack exchange,提问作者Cesar Pereira
相关产品推荐
相关产品推荐

