使用sparklyr/dplyr/pyspark统计用户组组合的事件发生与未发生次数
问题根因
你的代码运行结果不符合预期,核心是两点:
- 额外添加的
order_seq时间序列拆分逻辑和需求不匹配:需求是统计每个ID对应的全量所属组组合的事件触发情况,不需要按时间拆分同一个ID的多条记录 - 代码中引用了不存在的
group_a/group_b列,且没有先按ID维度做第一层聚合,得到每个ID的组属性和事件触发标记
R(sparklyr/dplyr通用实现)
代码逻辑完全兼容本地dplyr和sparklyr操作Spark DataFrame,无需修改:
# 第一步:按ID聚合,得到每个ID的组标记、是否触发过事件 id_agg <- input_data %>% group_by(id) %>% summarise( group_a = as.integer(any(group == "A")), group_b = as.integer(any(group == "B")), group_c = as.integer(any(group == "C")), has_event = as.integer(any(event == 1)) ) %>% ungroup() # 第二种输出格式(组名拼接)可以在这步加: # mutate(group_combo = paste0(case_when(group_a==1~"A"), case_when(group_b==1~",B"), case_when(group_c==1~",C")) %>% gsub("^,","",.)) # 第二步:按组组合聚合,统计数量 result <- id_agg %>% group_by(group_a, group_b, group_c) %>% summarise( event_occured = sum(has_event == 1), event_not_occured = sum(has_event == 0) ) %>% ungroup()
运行上述代码得到的结果和你要求的第一种输出格式完全一致。如果需要第二种拼接组名的格式,放开注释里的group_combo生成逻辑,后续按group_combo分组聚合即可。
PySpark实现
逻辑和R版本完全一致:
from pyspark.sql import functions as F # 第一步:按ID聚合 id_agg = input_data.groupBy("id").agg( F.max(F.when(F.col("group") == "A", 1).otherwise(0)).alias("group_a"), F.max(F.when(F.col("group") == "B", 1).otherwise(0)).alias("group_b"), F.max(F.when(F.col("group") == "C", 1).otherwise(0)).alias("group_c"), F.max(F.when(F.col("event") == 1, 1).otherwise(0)).alias("has_event") ) # 第二步:按组组合聚合 result = id_agg.groupBy("group_a", "group_b", "group_c").agg( F.sum(F.when(F.col("has_event") == 1, 1).otherwise(0)).alias("event_occured"), F.sum(F.when(F.col("has_event") == 0, 1).otherwise(0)).alias("event_not_occured") ) # 如果需要组名拼接格式,可以加如下逻辑: # id_agg = id_agg.withColumn("group_combo", # F.concat_ws(",", # F.when(F.col("group_a")==1, F.lit("A")), # F.when(F.col("group_b")==1, F.lit("B")), # F.when(F.col("group_c")==1, F.lit("C")) # ) # ) # 后续按group_combo分组聚合即可
内容的提问来源于stack exchange,提问作者Cyrus Mohammadian
相关产品推荐
相关产品推荐

