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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 15:36:26