PySpark中关联两个DataFrame并过滤已购买客户的推荐商品
解决PySpark DataFrame关联并标记购买状态的问题
核心逻辑
先把推荐记录(ps)的客户范围过滤成仅包含有购买记录(pp)的客户,再给每条推荐记录标记:如果该物品是客户已购买的,purchased设为1,否则设为0。
实现方案一(内关联+状态判断)
from pyspark.sql import functions as F # 按customer_id内关联,只保留有购买记录的客户的推荐数据 result_df = products_suggested.alias("ps") \ .join(products_purchased.alias("pp"), on="customer_id", how="inner") \ # 对比item_id,匹配则设1,否则0 .select( "ps.customer_id", "ps.item_id", F.when(F.col("ps.item_id") == F.col("pp.item_id"), 1).otherwise(0).alias("purchased") ) \ # 去重,避免ps中同一客户同一物品的重复推荐记录 .dropDuplicates(["customer_id", "item_id"])
实现方案二(左关联+过滤+coalesce)
如果习惯用左关联的方式,也可以先关联再过滤客户范围:
from pyspark.sql import functions as F result_df = products_suggested.alias("ps") \ # 按customer_id和item_id左关联 .join(products_purchased.alias("pp"), on=["customer_id", "item_id"], how="left") \ # 只保留有购买记录的客户(pp的customer_id不为空即代表该客户在pp中存在) .filter(F.col("pp.customer_id").isNotNull()) \ # 用coalesce取pp的purchased值(匹配到就是1),没匹配到就设0 .select( "ps.customer_id", "ps.item_id", F.coalesce(F.col("pp.purchased"), F.lit(0)).alias("purchased") )
关键说明
- 第一种方案用内关联直接限制了客户范围,效率更高;
- 第二种方案先左关联再过滤,逻辑更直观,适合理解关联关系的场景;
- 去重步骤根据ps的实际数据决定,如果ps本身没有重复的
customer_id+item_id,可以省略。
内容的提问来源于stack exchange,提问作者Andrew Ingalls
相关产品推荐
相关产品推荐

