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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 13:57:13