简介:这份资源面向推荐系统入门与进阶开发者,提供一套基于Spark MLlib实现的豆瓣电影推荐系统完整项目,帮助理解协同过滤在真实场景中的落地方式。项目以ALS算法为核心,覆盖数据预处理、训练测试集划分、参数调优、评分预测与RMSE、MAE等指标评估,并涉及覆盖率与多样性等推荐质量维度,适合作为大数据与人工智能方向的实战练习。压缩包共4个文件,约6.23MB,包含pom.xml依赖配置、Scala源码、Shell提交脚本及数据压缩包,结构紧凑,便于快速导入运行。目前已有595人学习下载。通过研读源码与数据,读者可掌握用户-物品交互建模、隐含特征向量求解及推荐结果生成流程,并理解如何结合物品相似度策略提升推荐多样性,为后续构建个性化推荐服务积累可复用的工程经验。
1. 豆瓣电影推荐系统:从 Spark ML 到可复现的离线推荐链路
豆瓣电影推荐系统这个题目,在人工智能大作业和毕设选题里出现的频率极高,但真正能跑通、能解释清楚每一行输出含义的并不多。我见过太多同学把 ALS 模型训练完,RMSE 打印出来就结束了,问他“给用户 u 推荐的前 10 部电影怎么来的”,答不上来。这篇笔记要解决的就是这个问题:用 Spark ML 的 ALS 算法,搭一条从豆瓣电影评分数据到 Top-N 推荐的完整离线链路,每一步都能复现,每个参数都能解释。
适合谁看?如果你正在做推荐系统相关的课程设计、毕设,或者刚转推荐方向想找一个能跑通的入门项目,这篇内容可以直接抄作业。如果你已经做过协同过滤,但说不清 implicitPrefs、冷启动、正则系数这些概念在实际数据上的表现,中间几章的参数分析和避坑记录会对你有用。整条链路基于 Spark 的 DataFrame 和 MLlib,不依赖深度学习框架,单机 8GB 内存就能跑通中等规模数据集。
2. Spark ML 的 ALS 到底在算什么:矩阵分解的直觉与选型理由
2.1 用户-物品评分矩阵为什么需要分解
推荐系统最原始的数据形态是一张巨大的稀疏矩阵:行是用户,列是电影,格子里是评分。豆瓣有数百万用户和数十万电影,但每个用户看过的电影通常只有几十到几百部,矩阵稀疏度往往超过 99%。这种矩阵直接做相似度计算,内存扛不住,而且大量缺失值让距离度量失去意义。
ALS(Alternating Least Squares,交替最小二乘)的思路是把这个大矩阵拆成两个小矩阵相乘:用户因子矩阵 U(用户数 × 隐因子数)和物品因子矩阵 V(电影数 × 隐因子数)。预测评分就是 U 的第 i 行和 V 的第 j 行做点积。隐因子数 k 通常取 10 到 200,相当于用 k 个潜在特征来描述一个用户或一部电影——可能是“偏文艺”“爱看动作”“对老片容忍度高”这类无法直接命名但数值上有效的维度。
选 ALS 而不是基于邻域的方法,核心理由有三条:第一,Spark ML 的 ALS 实现是分布式的,能处理单机放不下的评分数据;第二,它天然支持隐式反馈(implicitPrefs),豆瓣的“看过”行为可以转化为置信度而不是显式评分;第三,交替求解的过程可以并行化,每轮固定一边求另一边,是闭式解,收敛行为比随机梯度下降更可控。
2.2 显式反馈与隐式反馈在豆瓣数据上的取舍
豆瓣数据有两种可用信号:显式评分(1 到 5 星)和隐式行为(看过、想看、评论)。显式评分最直接,但问题是稀疏且存在用户偏置——有人习惯打 3 星,有人动不动就 5 星。隐式反馈把“看过”当作正例,把“没看过”当作弱负例,用置信度加权,通常在实际系统中效果更稳。
Spark ML 的 ALS 通过implicitPrefs参数切换两种模式。设为true时,评分值被解释为置信度,算法优化的是偏好排序而不是评分误差;设为false时,直接最小化预测评分和真实评分的平方误差。我的经验是:如果数据里显式评分覆盖率超过 5%,先用显式模式跑基线;如果评分极稀疏但行为日志丰富,切隐式模式。豆瓣公开数据集通常显式评分就够用,所以下面以显式模式为主线,隐式模式的参数差异在避坑章节展开。
2.3 最小可跑通的 ALS 训练代码
先看一段能在本地 Spark 环境直接跑的最小代码,数据格式是userId, movieId, rating, timestamp的 CSV。
from pyspark.sql import SparkSession from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator # 初始化 SparkSession,单机模式用 local[*] 吃满 CPU spark = SparkSession.builder \ .appName("DoubanMovieALS") \ .master("local[*]") \ .config("spark.driver.memory", "4g") \ .getOrCreate() # 读取评分数据,显式指定 schema 避免类型推断翻车 ratings = spark.read.csv( "data/ratings.csv", header=True, inferSchema=True ).select("userId", "movieId", "rating") # 按 8:2 切训练集和测试集,固定种子保证可复现 train, test = ratings.randomSplit([0.8, 0.2], seed=42) # 定义 ALS 模型,核心参数先给一组经验值 als = ALS( userCol="userId", itemCol="movieId", ratingCol="rating", rank=50, # 隐因子维度 maxIter=10, # 交替迭代轮数 regParam=0.1, # 正则化系数,防过拟合 implicitPrefs=False, # 显式评分模式 coldStartStrategy="drop", # 预测时丢弃冷启动用户/物品 nonnegative=True, # 因子非负,提升可解释性 seed=42 ) # 训练 model = als.fit(train) # 在测试集上预测并评估 predictions = model.transform(test) evaluator = RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ) rmse = evaluator.evaluate(predictions) print(f"Test RMSE = {rmse:.4f}") # 给每个用户生成 Top-10 推荐 user_recs = model.recommendForAllUsers(10) user_recs.show(5, truncate=False)这段代码的逻辑链条是:读数据 → 切分 → 定义 ALS → 训练 → 评估 → 生成推荐。几个关键点需要展开。rank=50是隐因子数,太小欠拟合,太大过拟合且训练慢,50 是中等规模数据集的常用起点。maxIter=10通常够收敛,但如果你发现 RMSE 还在明显下降,可以加到 15 或 20。regParam=0.1控制正则强度,值越大模型越保守,对稀疏数据的过拟合抑制越明显。coldStartStrategy="drop"很重要——测试集里可能出现训练集没见过的用户或电影,不丢弃的话预测结果是 NaN,RMSE 直接变 NaN,这是新手最常见的翻车点之一。
nonnegative=True让分解出的因子非负,好处是推荐结果更容易解释(因子可以理解为“正向偏好强度”),代价是可能略微抬高 RMSE。如果你的目标只是排序质量,可以关掉;如果要做可解释推荐,建议打开。
3. 豆瓣数据从原始 CSV 到 ALS 输入:清洗、编码与特征工程
3.1 豆瓣评分数据的典型脏法
从豆瓣抓取或从公开数据集拿到的评分数据,通常有这几类问题:用户 ID 和电影 ID 是字符串(比如u12345、tt0111161),ALS 要求整数索引;评分有缺失或超出 1 到 5 的范围;同一用户对同一电影有多条记录(重复评分);时间戳格式不统一。不做清洗直接喂给 ALS,轻则报类型错误,重则训练出的模型完全不可用。
清洗的目标是得到一张干净的userId: Int, movieId: Int, rating: Float三元组表。注意 ALS 在 Spark ML 里要求用户列和物品列是整数类型,评分列是浮点类型。字符串 ID 必须做索引编码,而且编码要稳定——训练集和测试集必须用同一套映射,否则同一个用户在两边被编成不同整数,模型直接错乱。
3.2 用 StringIndexer 做 ID 编码的完整步骤
from pyspark.ml.feature import StringIndexer from pyspark.sql.functions import col, when, count, desc # 假设原始数据列名是 user_id, movie_id, score raw = spark.read.csv("data/douban_raw.csv", header=True, inferSchema=True) # 1. 过滤评分范围,只保留 1-5 的整数评分 raw = raw.filter((col("score") >= 1) & (col("score") <= 5)) # 2. 去重:同一用户对同一电影保留最新一条 from pyspark.sql.window import Window from pyspark.sql.functions import row_number window = Window.partitionBy("user_id", "movie_id").orderBy(desc("timestamp")) raw = raw.withColumn("rn", row_number().over(window)) \ .filter(col("rn") == 1).drop("rn") # 3. 字符串 ID 转整数索引 user_indexer = StringIndexer(inputCol="user_id", outputCol="userId", handleInvalid="skip") movie_indexer = StringIndexer(inputCol="movie_id", outputCol="movieId", handleInvalid="skip") user_indexer_model = user_indexer.fit(raw) movie_indexer_model = movie_indexer.fit(raw) indexed = user_indexer_model.transform(raw) indexed = movie_indexer_model.transform(indexed) # 4. 选出 ALS 需要的三列,评分转 float ratings = indexed.select( col("userId").cast("int"), col("movieId").cast("int"), col("score").cast("float").alias("rating") ) # 5. 检查基本统计量 print(f"用户数: {ratings.select('userId').distinct().count()}") print(f"电影数: {ratings.select('movieId').distinct().count()}") print(f"评分数: {ratings.count()}") ratings.groupBy("rating").count().orderBy("rating").show()这段代码里有两个容易忽略的点。第一,StringIndexer的handleInvalid="skip"表示遇到新类别时跳过而不是报错,但在训练集上 fit 之后,测试集 transform 时如果出现训练集没有的 ID,这些行会被跳过——这是合理的,因为 ALS 本来也处理不了冷启动物品。第二,去重逻辑用窗口函数按时间戳取最新一条,比简单dropDuplicates更符合业务含义:用户改了评分,应该以最后一次为准。
编码完成后,userId和movieId都是从 0 开始的连续整数,但注意StringIndexer默认按出现频率降序编码,所以编号本身没有大小含义,不要拿去做数值比较。
3.3 评分分布检查与偏置处理
清洗完一定要看一眼评分分布。豆瓣用户打分普遍偏高,3 星以下很少,这会导致模型倾向于预测高分。如果发现 4 星和 5 星占了 80% 以上,有两个处理方向:一是对评分做中心化,减去用户均值或全局均值;二是改用隐式反馈模式,把评分高低转化为置信度。
Spark ML 的 ALS 本身没有内置中心化,需要手动做。简单做法是计算全局均值,训练时用rating - global_mean,预测后再加回来。更精细的做法是按用户去均值,但那样需要额外维护用户偏置表,工程复杂度上升。对于课程设计级别的项目,全局均值中心化通常够用,RMSE 能降 0.05 到 0.1。
from pyspark.sql.functions import avg, stddev stats = ratings.select( avg("rating").alias("mean"), stddev("rating").alias("std") ).collect()[0] print(f"评分均值: {stats['mean']:.3f}, 标准差: {stats['std']:.3f}") # 全局均值中心化 global_mean = stats["mean"] ratings_centered = ratings.withColumn( "rating", col("rating") - global_mean )中心化之后,ALS 的ratingCol换成rating(已经是中心化后的值),预测时记得把global_mean加回去再算 RMSE,否则评估指标没有意义。
4. ALS 参数调优:rank、regParam、maxIter 怎么定
4.1 用 CrossValidator 做网格搜索的代价
Spark ML 提供了CrossValidator和ParamGridBuilder,理论上可以自动搜参。但 ALS 的训练成本随 rank 和 maxIter 线性增长,三折交叉验证乘以参数组合数,单机跑一天都跑不完。我的做法是:先用小规模采样数据(比如 10% 用户)做粗筛,确定参数大致范围,再在全量数据上精调一两个关键参数。
粗筛阶段可以固定 maxIter=5,只搜 rank 和 regParam。rank 候选 [10, 30, 50, 100],regParam 候选 [0.01, 0.05, 0.1, 0.5]。每组跑完记录 RMSE,画一张热力图,通常能看到一个明显的低谷区域。
4.2 三个核心参数的交互影响
rank 决定模型容量。太小(比如 5)时,用户和电影被压缩到极低维空间,区分度不够,RMSE 偏高;太大(比如 200)时,每个因子分到的数据变少,过拟合风险上升,而且训练时间显著增加。在豆瓣中等规模数据上,rank 在 30 到 80 之间通常能找到较优值。
regParam 控制正则化强度。它的作用和 rank 相反:rank 大时,需要更大的 regParam 来抑制过拟合;rank 小时,regParam 可以小一些。两者要联合调,单独调一个往往得不到最优。
maxIter 是迭代轮数。ALS 每轮都有闭式解,收敛通常较快。观察训练日志里的 RMSE 变化,如果连续三轮下降幅度小于 0.001,就可以停了。盲目设 50 轮除了浪费时间,还可能因为过拟合导致测试集 RMSE 反弹。
from pyspark.ml.tuning import ParamGridBuilder, CrossValidator # 小规模采样做粗筛 sample_users = ratings.select("userId").distinct().sample(False, 0.1, seed=42) sample_ratings = ratings.join(sample_users, on="userId") sample_train, sample_test = sample_ratings.randomSplit([0.8, 0.2], seed=42) als_tune = ALS( userCol="userId", itemCol="movieId", ratingCol="rating", maxIter=5, coldStartStrategy="drop", seed=42 ) param_grid = ParamGridBuilder() \ .addGrid(als_tune.rank, [10, 30, 50, 100]) \ .addGrid(als_tune.regParam, [0.01, 0.05, 0.1, 0.5]) \ .build() evaluator = RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ) cv = CrossValidator( estimator=als_tune, estimatorParamMaps=param_grid, evaluator=evaluator, numFolds=3, parallelism=2, # 并行跑 2 个模型,吃内存 seed=42 ) cv_model = cv.fit(sample_train) best_rank = cv_model.bestModel.rank best_reg = cv_model.bestModel._java_obj.parent().getRegParam() print(f"粗筛最优: rank={best_rank}, regParam={best_reg}")parallelism=2表示同时训练两个参数组合,能加速但吃内存。如果机器内存小于 8GB,建议设为 1,否则容易 OOM。粗筛得到最优参数后,用全量训练集重新训练,maxIter 可以适当加大到 10 到 15。
4.3 评估指标不只看 RMSE
RMSE 衡量的是评分预测误差,但推荐系统的核心目标是排序质量。一个 RMSE 很低的模型,可能只是学会了预测用户已经看过的电影的高分,对新电影的排序能力未必好。所以除了 RMSE,还应该看 Precision@K、Recall@K 或 NDCG@K。
Spark ML 没有内置这些排序指标,需要自己实现。一个简化做法是:对测试集中每个用户,取模型预测分数最高的 K 部电影,看有多少部真的出现在该用户的测试集正例中,算命中率。虽然粗糙,但比只看 RMSE 更能反映推荐效果。
# 给测试集用户生成 Top-10 推荐 test_users = test.select("userId").distinct() recs = model.recommendForUserSubset(test_users, 10) # 展开推荐结果,和测试集正例(评分>=4)做命中统计 from pyspark.sql.functions import explode recs_exploded = recs.select("userId", explode("recommendations").alias("rec")) \ .select("userId", col("rec.movieId").alias("movieId")) test_positive = test.filter(col("rating") >= 4).select("userId", "movieId") hits = recs_exploded.join(test_positive, on=["userId", "movieId"], how="inner") hit_count = hits.count() total_recs = recs_exploded.count() print(f"命中率: {hit_count / total_recs:.4f}")这个命中率不是标准 Precision@K,因为分母是推荐总数而不是用户数乘以 K,但作为快速对比不同参数的相对指标够用了。
5. 避坑与排查:ALS 训练和推荐环节的 5 个血泪教训
5.1 预测结果出现 NaN,RMSE 直接变 NaN
现象:训练完模型,transform(test)之后评估,RMSE 打印出来是nan。
原因:测试集里存在训练集没出现过的 userId 或 movieId,ALS 对冷启动用户/物品的预测默认返回 NaN。如果不处理,RegressionEvaluator算出来的就是 NaN。
解决:定义 ALS 时加coldStartStrategy="drop",预测阶段自动丢弃含 NaN 的行。如果不想丢数据,可以改用coldStartStrategy="nan"然后手动填充全局均值,但推荐质量会下降。更根本的做法是在切分数据时保证测试集的用户和物品都出现在训练集中,可以用randomSplit后做一次交集过滤。
5.2 显式评分模式下评分未做浮点转换导致类型报错
现象:als.fit(train)报IllegalArgumentException: requirement failed: Column rating must be of type float but was actually int。
原因:CSV 读进来时inferSchema=True可能把评分推断成整数,而 ALS 要求ratingCol是浮点类型。
解决:在读数据后显式cast("float"),或者在 schema 里直接指定FloatType()。这个错误信息其实很明确,但新手容易忽略,因为报错发生在 fit 阶段而不是读数据阶段。
5.3 rank 设得太大导致单机 OOM
现象:训练到一半抛OutOfMemoryError: Java heap space,或者 Spark 任务卡在某个 stage 不动。
原因:rank 增大时,用户因子矩阵和物品因子矩阵的维度线性增长,加上 ALS 每轮迭代要缓存中间结果,内存占用是 rank 的倍数。单机 8GB 内存跑 rank=200 加 maxIter=20,很容易撑爆。
解决:先降 rank 到 50 以下,或者减小 maxIter。如果必须用大 rank,可以调大spark.driver.memory和spark.executor.memory,但单机有上限。另一个方向是减少数据量,比如只保留评分次数超过 5 次的用户和电影,稀疏度降低后内存压力也会小很多。
5.4 推荐结果全是热门电影,长尾物品出不来
现象:给不同用户生成的 Top-10 推荐高度重合,翻来覆去就是那几部高分经典片。
原因:ALS 在显式评分上优化的是评分预测误差,热门电影评分多、均值高,因子向量被训练得偏向全局高分方向,导致对所有用户都推荐类似的片子。这是协同过滤的经典问题,不是 bug。
解决:三个方向。一是改用隐式反馈,把“看过”作为正例,热门电影的正例多但置信度可以按流行度打折;二是对物品因子做流行度惩罚,推荐分数减去一个和物品流行度正相关的项;三是在训练数据里对热门电影降采样,减少它们在损失函数中的权重。课程设计级别至少要做第一个或第三个,否则推荐结果没有说服力。
5.5 训练集和测试集编码不一致导致用户错位
现象:模型训练时 RMSE 正常,但推荐结果明显不合理,比如给只看过动画片的用户推荐恐怖片。
原因:如果训练集和测试集分别做StringIndexer,同一个字符串 ID 可能被编成不同整数。更隐蔽的情况是,先切分再编码,训练集 fit 的 indexer 没有应用到测试集,两边编码体系不一致。
解决:编码必须在切分之前做,或者用训练集 fit 出的 indexer 模型去 transform 测试集。正确顺序是:原始数据 → 清洗 → 编码 → 切分 → 训练/测试。切分之后不要再做任何改变 ID 映射的操作。
6. 从离线推荐到可解释输出:让 ALS 结果能讲清楚
6.1 用物品因子做相似电影检索
ALS 训练完,物品因子矩阵model.itemFactors是一个 DataFrame,每行是id和features(一个数组)。两部电影的相似度可以用因子向量的余弦相似度衡量。这个能力可以用来做“看了又看”或“相似推荐”,也是验证模型是否学到有意义结构的好方法。
from pyspark.sql.functions import udf, array, col from pyspark.ml.linalg import Vectors, VectorUDT import numpy as np # 取出物品因子 item_factors = model.itemFactors item_factors.show(3, truncate=False) # 定义余弦相似度 UDF def cosine_sim(v1, v2): a = np.array(v1) b = np.array(v2) return float(np.dot(a, b) / (np.linalg.norm(a) * np.linalg.norm(b) + 1e-10)) # 以某部电影为例,找最相似的 10 部 target_id = 50 # 假设这是某部电影的编码 ID target_vec = item_factors.filter(col("id") == target_id).select("features").collect()[0][0] # 广播目标向量,计算所有物品的相似度 item_factors_local = item_factors.collect() sims = [] for row in item_factors_local: if row["id"] != target_id: sim = cosine_sim(target_vec, row["features"]) sims.append((row["id"], sim)) sims.sort(key=lambda x: x[1], reverse=True) print("最相似的 10 部电影 ID:", [s[0] for s in sims[:10]])这段代码在单机上跑没问题,但如果物品数量到几十万,collect()会把所有因子拉到 driver,内存吃紧。生产环境应该用 Spark 的分布式矩阵运算,或者用ColumnSimilarities做近似计算。课程设计级别,采样几千部电影做相似度检索就够了。
6.2 推荐理由的生成思路
可解释推荐是现在人工智能应用里越来越被看重的能力。ALS 本身不输出“为什么推荐”,但可以从因子向量反推。一个简单做法是:对推荐给用户的每部电影,找到该用户已评分最高的几部电影,计算它们和推荐电影在因子空间中的相似度,把相似度最高的那部作为“因为你看了 X”的理由。
# 伪代码思路,实际实现需要 join 用户历史高分电影和推荐结果 # 1. 取用户 u 评分 >= 4 的电影集合 H # 2. 取推荐给 u 的电影集合 R # 3. 对 R 中每部电影 r,在 H 中找因子相似度最高的 h # 4. 输出 "因为你看了 h,推荐 r"这个思路的局限是,因子空间的相似度不等于内容相似度,有时候理由看起来会有点玄学。但作为课程设计的加分项,能跑通并展示出来,已经比只打印 RMSE 强很多。
6.3 我踩过的一个坑:不要用测试集调参
最后说一个我自己的教训。早期做这个项目时,我拿测试集 RMSE 来选 rank 和 regParam,看到某个参数组合测试集 RMSE 最低就直接用了。后来才意识到,这等于用测试集做了模型选择,评估结果偏乐观,实际部署时效果会打折扣。正确做法是切出验证集,用验证集调参,测试集只在最后评估一次。如果数据量实在不够,至少用交叉验证,不要反复在同一个测试集上试参数。
另一个习惯是:每次训练完保存模型和参数配置,包括 rank、regParam、maxIter、训练集大小、RMSE。过两周回头看,没有这些记录根本记不清哪个模型对应哪组参数。Spark ML 的model.save()可以保存模型,但参数配置要自己写到日志或配置文件里。
希望帮到你。
本文还有配套的精品资源,点击获取