PySpark内连接引发笛卡尔积问题:ALS电影推荐结果处理异常
这个问题我之前也踩过坑!核心原因很简单:ALS输出的movie_id和rating数组是一一对应的(第一个movie对应第一个评分,以此类推),但如果你把两个数组分开拆成独立DataFrame再只按user_id连接,PySpark不知道元素的位置关系,就会把所有movie_id和所有rating两两配对,自然就产生了笛卡尔积。
错误做法的问题所在
举个例子,假设你的user_recs_one数据是这样的:
+-------+----------------+---------------+
|user_id|movie_id |rating |
+-------+----------------+---------------+
|1 |[123, 456, 789] |[4.2, 3.8, 4.5]|
+-------+----------------+---------------+
如果像你之前那样拆分后连接:
# 错误示例:拆分后仅按user_id连接 movies_df = user_recs_one.select("user_id", explode("movie_id").alias("movie_id")) ratings_df = user_recs_one.select("user_id", explode("rating").alias("rating")) joined_df = movies_df.join(ratings_df, on="user_id")
你会得到9行数据(3个movie × 3个rating),完全不是预期的一一对应结果。
两种正确的解决方法
方法1:用arrays_zip打包数组后拆分(推荐,更简洁)
先把movie_id和rating数组打包成结构体数组,这样每个结构体里的movie和rating是对应的,再拆分这个结构体数组就不会乱序:
from pyspark.sql.functions import arrays_zip, explode, col # 第一步:将两个数组打包成结构体数组 zipped_recs = user_recs_one.select( "user_id", arrays_zip("movie_id", "rating").alias("recommendation_pairs") ) # 第二步:拆分结构体数组,提取对应的movie_id和rating final_recs = zipped_recs.select( "user_id", explode("recommendation_pairs").alias("pair") ).select( "user_id", col("pair.movie_id").alias("movie_id"), col("pair.rating").alias("rating") ) final_recs.show(truncate=False)
执行后会得到3行正确的对应结果:
+-------+--------+------+
|user_id|movie_id|rating|
+-------+--------+------+
|1 |123 |4.2 |
|1 |456 |3.8 |
|1 |789 |4.5 |
+-------+--------+------+
方法2:用posexplode获取位置索引后连接
如果需要保留元素的位置信息(比如后续要按推荐顺序排序),可以用posexplode拆分数组时同时获取元素的位置索引,然后通过user_id + 位置索引来连接,保证对应位置的元素匹配:
from pyspark.sql.functions import posexplode # 拆分movie_id,保留位置pos movies_df = user_recs_one.select( "user_id", posexplode("movie_id").alias("pos", "movie_id") ) # 拆分rating,保留相同的位置pos ratings_df = user_recs_one.select( "user_id", posexplode("rating").alias("pos", "rating") ) # 按user_id和pos同时连接,最后删除pos列 final_recs = movies_df.join(ratings_df, on=["user_id", "pos"], how="inner").drop("pos") final_recs.show(truncate=False)
这个方法也能得到和上面一样的正确结果,适合需要处理位置相关逻辑的场景。
关键提醒
ALS的推荐结果中,movie_id和rating数组的元素顺序是严格对应的,所以拆分时一定要保持元素的位置关联,不能单独拆分后只按user_id连接,否则必然会出现笛卡尔积问题。
内容的提问来源于stack exchange,提问作者Anindya Saha

