PySpark聚合DataFrame生成sent/open/click多列统计实现方法
PySpark 按客户维度统计邮件营销行为指标实现方案
核心注意点
- 原代码存在拼写错误:
groupfBy应为groupBy - 分组聚合时需要把同粒度的维度字段(mkt_channel_name、mkt_channel_category)都放入分组字段,否则聚合后会丢失这些列
- 条件计数使用
F.when()搭配聚合函数实现,不需要写多次分组逻辑
方案1:聚合去重(每个客户+渠道仅返回1行统计结果)
如果不需要保留原始明细行的mkt_channel_subcategory字段,直接用分组聚合实现,代码如下:
from pyspark.sql import functions as F result_df = d.groupBy( "CUSTOMER_ID", "mkt_channel_id", "mkt_channel_name", "mkt_channel_category" ).agg( # sent统计规则:subcategory为send/delivery均计入 F.count(F.when(F.col("mkt_channel_subcategory").isin("send", "delivery"), 1)).alias("sent"), F.count(F.when(F.col("mkt_channel_subcategory") == "open", 1)).alias("open"), F.count(F.when(F.col("mkt_channel_subcategory") == "click", 1)).alias("click") )
方案2:保留明细行(每行都附带同维度统计值)
如果你需要保留原始所有明细列(包括mkt_channel_subcategory),让每一行明细都带上对应客户+渠道的三个统计值,使用窗口函数实现:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义统计窗口:按客户+营销渠道分区 win_spec = Window.partitionBy("CUSTOMER_ID", "mkt_channel_id") result_df = d.withColumn( "sent", F.count(F.when(F.col("mkt_channel_subcategory").isin("send", "delivery"), 1)).over(win_spec) ).withColumn( "open", F.count(F.when(F.col("mkt_channel_subcategory") == "open", 1)).over(win_spec) ).withColumn( "click", F.count(F.when(F.col("mkt_channel_subcategory") == "click", 1)).over(win_spec) )
两种方案的区别:方案1返回的是聚合后的宽表,数据量更小;方案2不改变原始数据行数,保留所有明细,适合需要明细+统计值同时展示的场景,可根据业务需求选择。
内容的提问来源于stack exchange,提问作者merkle
相关产品推荐
相关产品推荐

