PySpark如何基于宠物明细表生成共现计数聚合表
PySpark 实现用户饲养宠物共现计数矩阵
问题场景
已通过explode_outer将原表嵌套数组格式的pets字段拆解为扁平化明细表,表结构为(user: string, pet: string),每行对应用户饲养的单只宠物。需要生成的共现计数矩阵满足:
- 行、列索引均为全量去重后的宠物类型
- 单元格值为对应两种宠物被同一用户共同饲养的次数
- 同宠物对应的对角线单元格值固定为0
- 可直接支撑“查询指定宠物适配度最高的关联宠物”类分析需求
实现思路
核心逻辑是同用户的宠物列表做自连接配对,过滤同宠物配对后分组计数,最后将长表转换为宽表矩阵即可,步骤拆解:
- 对扁平化明细表按
user分组,收集每个用户对应的去重宠物列表,避免单用户重复饲养同一种宠物导致计数虚高 - 对收集到的宠物列表做自交叉配对,生成同一用户下所有两两宠物的组合
- 过滤掉宠物A与宠物B相同的配对,按
(pet_a, pet_b)维度分组统计用户数,得到共现计数长表 - 以pet_a为行、pet_b为列、共现次数为值做透视,空值填充为0,对角线值统一重置为0,得到最终矩阵
代码实现
假设拆解完成的扁平化明细表名为flat_pet_df,包含user、pet两个字段,完整实现代码如下:
from pyspark.sql import functions as F # 按用户聚合去重宠物列表 user_pet_list_df = flat_pet_df.groupBy("user")\ .agg(F.collect_set("pet").alias("pet_list")) # 自连接生成同用户下所有宠物两两配对 pet_pairs_df = user_pet_list_df.select( "user", F.explode(F.col("pet_list")).alias("pet_a") ).join( user_pet_list_df.select("user", F.explode(F.col("pet_list")).alias("pet_b")), on="user", how="inner" ) # 过滤同宠物配对,统计共现用户数 cooccur_long_df = pet_pairs_df.filter(F.col("pet_a") != F.col("pet_b"))\ .groupBy("pet_a", "pet_b")\ .agg(F.countDistinct("user").alias("cooccur_cnt")) # 获取全量宠物类型列表,用于透视保证列完整 all_pets = [row.pet for row in flat_pet_df.select("pet").distinct().collect()] # 透视生成宽表矩阵,空值填0 cooccur_matrix_df = cooccur_long_df.groupBy("pet_a")\ .pivot("pet_b", all_pets)\ .agg(F.first("cooccur_cnt"))\ .na.fill(0) # 兜底将对角线值统一置为0 for pet in all_pets: cooccur_matrix_df = cooccur_matrix_df.withColumn( pet, F.when(F.col("pet_a") == F.lit(pet), 0).otherwise(F.col(pet)) )
结果使用示例
查询和狗(dog)适配度最高的前3种宠物,直接取对应列降序排列即可:
cooccur_matrix_df.select("pet_a", "dog")\ .orderBy(F.col("dog").desc())\ .limit(3)\ .show()
效果验证:如果蜘蛛仅在用户bb处与狗共同饲养,最终矩阵中spider行dog列的计数值为1,符合预期。
内容的提问来源于stack exchange,提问作者Hyunuk
相关产品推荐
相关产品推荐

