PySpark如何遍历匹配两个DataFrame计算推荐系统命中率
PySpark实现推荐系统命中率计算
原有代码的问题
- 语法不符合PySpark规范:PySpark DataFrame不支持Python原生的逐行for循环遍历写法;判断相等误用赋值符
=,应该用== - 统计逻辑错误:未做用户维度去重,单用户多次命中会被重复计数;分母用
recent_df.count()统计的是总购买记录条数,不是去重后的有效用户总数 - 性能问题:双层循环本质是笛卡尔积遍历,数据量稍大就会出现计算卡顿甚至任务失败
正确实现代码
核心思路是通过DataFrame原生join操作匹配命中记录,再做用户维度去重统计,完全适配Spark分布式计算特性:
from pyspark.sql import functions as F def calc_hit_ratio(recommendation_df, recent_df): # 统计基准:recent_df中去重后的总用户数 total_user = recent_df.select("new_party_id").distinct().count() # 内连接匹配:用户ID一致 + 购买商品在推荐列表中 match_df = recent_df.join( recommendation_df, on=( (recent_df.new_party_id == recommendation_df.party_id) & (recent_df.merch_store_code == recommendation_df.merch_store_code) ), how="inner" ) # 命中用户去重:单用户多条命中记录只算1次 hit_user = match_df.select("new_party_id").distinct().count() # 计算最终命中率 return hit_user / total_user # 调用示例 hit_ratio = calc_hit_ratio(recommendation_df, recent_df) print(f"推荐命中率为:{hit_ratio:.4f}")
逻辑校验点
- 匹配时严格对齐用户ID:
recommendation_df的用户字段为party_id,recent_df的用户字段为new_party_id,不会出现跨用户匹配商品的错误 - 单用户只要存在任意一条购买商品命中推荐列表即计数1次,符合规则要求
- 所有统计基于Spark原生算子,无逐行遍历逻辑,大数据量下计算效率稳定
内容的提问来源于stack exchange,提问作者jcng2308
相关产品推荐
相关产品推荐

