如何通过PySpark实现按ID统计有效邮箱占比?
PySpark 按ID统计有效邮箱占比解决方案
你的代码问题分析
- 窗口函数分区错误:你把
email、name、surname、validity都加入了分区,这会导致每个ID下的不同记录被拆分成多个小分组,row_number()根本无法实现按ID聚合的目的。 - 有效数计算逻辑错误:
withColumn里嵌套df.select(...).count()的写法完全不符合PySpark语法,这种方式无法将聚合结果映射回原数据的每一行。
正确实现方法
要实现按ID统计有效邮箱占比并添加到原数据中,用窗口函数是最直接的方式,无需额外join操作:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义仅按ID分区的窗口 id_window = Window.partitionBy("ID") # 计算每个ID的有效邮箱数、总邮箱数,再算出占比 df_result = df.withColumn( "total_emails_per_id", F.count("*").over(id_window) ).withColumn( "valid_emails_per_id", F.sum(F.when(F.col("Validity") == "valid", 1).otherwise(0)).over(id_window) ).withColumn( "valid_ratio", F.round(F.col("valid_emails_per_id") / F.col("total_emails_per_id"), 2) # 保留两位小数 )
代码说明
- 窗口仅按
ID分区,确保所有同ID的记录共享聚合结果 count("*").over(id_window):统计每个ID下的总记录数sum(when(...)):统计每个ID下Validity为valid的记录数- 最后通过除法计算占比,用
round控制小数位数
如果你的数据里存在重复的邮箱记录(同一个ID+邮箱多次出现),需要先去重再统计:
# 先按ID和邮箱去重,避免重复统计 df_distinct = df.dropDuplicates(["ID", "email"]) # 再用上面的窗口函数逻辑计算占比 id_window = Window.partitionBy("ID") df_result = df_distinct.withColumn( "total_emails_per_id", F.count("*").over(id_window) ).withColumn( "valid_emails_per_id", F.sum(F.when(F.col("Validity") == "valid", 1).otherwise(0)).over(id_window) ).withColumn( "valid_ratio", F.round(F.col("valid_emails_per_id") / F.col("total_emails_per_id"), 2) ) # 如果需要关联回原数据(保留所有原始记录),可以用join df_final = df.join(df_result.select("ID", "valid_ratio"), on="ID", how="left")
内容的提问来源于stack exchange,提问作者indritkalaj
相关产品推荐
相关产品推荐

