You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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)  # 保留两位小数
)

代码说明

  1. 窗口仅按ID分区,确保所有同ID的记录共享聚合结果
  2. count("*").over(id_window):统计每个ID下的总记录数
  3. sum(when(...)):统计每个ID下Validity为valid的记录数
  4. 最后通过除法计算占比,用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 13:15:58