ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于Spark ML的豆瓣电影推荐系统:ALS算法实战与调优

基于Spark ML的豆瓣电影推荐系统:ALS算法实战与调优 简介这份资源面向推荐系统入门与进阶开发者提供一套基于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%。这种矩阵直接做相似度计算内存扛不住而且大量缺失值让距离度量失去意义。ALSAlternating 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, headerTrue, inferSchemaTrue ).select(userId, movieId, rating) # 按 8:2 切训练集和测试集固定种子保证可复现 train, test ratings.randomSplit([0.8, 0.2], seed42) # 定义 ALS 模型核心参数先给一组经验值 als ALS( userColuserId, itemColmovieId, ratingColrating, rank50, # 隐因子维度 maxIter10, # 交替迭代轮数 regParam0.1, # 正则化系数防过拟合 implicitPrefsFalse, # 显式评分模式 coldStartStrategydrop, # 预测时丢弃冷启动用户/物品 nonnegativeTrue, # 因子非负提升可解释性 seed42 ) # 训练 model als.fit(train) # 在测试集上预测并评估 predictions model.transform(test) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(predictions) print(fTest RMSE {rmse:.4f}) # 给每个用户生成 Top-10 推荐 user_recs model.recommendForAllUsers(10) user_recs.show(5, truncateFalse)这段代码的逻辑链条是读数据 → 切分 → 定义 ALS → 训练 → 评估 → 生成推荐。几个关键点需要展开。rank50是隐因子数太小欠拟合太大过拟合且训练慢50 是中等规模数据集的常用起点。maxIter10通常够收敛但如果你发现 RMSE 还在明显下降可以加到 15 或 20。regParam0.1控制正则强度值越大模型越保守对稀疏数据的过拟合抑制越明显。coldStartStrategydrop很重要——测试集里可能出现训练集没见过的用户或电影不丢弃的话预测结果是 NaNRMSE 直接变 NaN这是新手最常见的翻车点之一。nonnegativeTrue让分解出的因子非负好处是推荐结果更容易解释因子可以理解为“正向偏好强度”代价是可能略微抬高 RMSE。如果你的目标只是排序质量可以关掉如果要做可解释推荐建议打开。3. 豆瓣数据从原始 CSV 到 ALS 输入清洗、编码与特征工程3.1 豆瓣评分数据的典型脏法从豆瓣抓取或从公开数据集拿到的评分数据通常有这几类问题用户 ID 和电影 ID 是字符串比如u12345、tt0111161ALS 要求整数索引评分有缺失或超出 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, headerTrue, inferSchemaTrue) # 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(inputColuser_id, outputColuserId, handleInvalidskip) movie_indexer StringIndexer(inputColmovie_id, outputColmovieId, handleInvalidskip) 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的handleInvalidskip表示遇到新类别时跳过而不是报错但在训练集上 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% 用户做粗筛确定参数大致范围再在全量数据上精调一两个关键参数。粗筛阶段可以固定 maxIter5只搜 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, seed42) sample_ratings ratings.join(sample_users, onuserId) sample_train, sample_test sample_ratings.randomSplit([0.8, 0.2], seed42) als_tune ALS( userColuserId, itemColmovieId, ratingColrating, maxIter5, coldStartStrategydrop, seed42 ) 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( metricNamermse, labelColrating, predictionColprediction ) cv CrossValidator( estimatorals_tune, estimatorParamMapsparam_grid, evaluatorevaluator, numFolds3, parallelism2, # 并行跑 2 个模型吃内存 seed42 ) 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})parallelism2表示同时训练两个参数组合能加速但吃内存。如果机器内存小于 8GB建议设为 1否则容易 OOM。粗筛得到最优参数后用全量训练集重新训练maxIter 可以适当加大到 10 到 15。4.3 评估指标不只看 RMSERMSE 衡量的是评分预测误差但推荐系统的核心目标是排序质量。一个 RMSE 很低的模型可能只是学会了预测用户已经看过的电影的高分对新电影的排序能力未必好。所以除了 RMSE还应该看 PrecisionK、RecallK 或 NDCGK。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], howinner) hit_count hits.count() total_recs recs_exploded.count() print(f命中率: {hit_count / total_recs:.4f})这个命中率不是标准 PrecisionK因为分母是推荐总数而不是用户数乘以 K但作为快速对比不同参数的相对指标够用了。5. 避坑与排查ALS 训练和推荐环节的 5 个血泪教训5.1 预测结果出现 NaNRMSE 直接变 NaN现象训练完模型transform(test)之后评估RMSE 打印出来是nan。原因测试集里存在训练集没出现过的 userId 或 movieIdALS 对冷启动用户/物品的预测默认返回 NaN。如果不处理RegressionEvaluator算出来的就是 NaN。解决定义 ALS 时加coldStartStrategydrop预测阶段自动丢弃含 NaN 的行。如果不想丢数据可以改用coldStartStrategynan然后手动填充全局均值但推荐质量会下降。更根本的做法是在切分数据时保证测试集的用户和物品都出现在训练集中可以用randomSplit后做一次交集过滤。5.2 显式评分模式下评分未做浮点转换导致类型报错现象als.fit(train)报IllegalArgumentException: requirement failed: Column rating must be of type float but was actually int。原因CSV 读进来时inferSchemaTrue可能把评分推断成整数而 ALS 要求ratingCol是浮点类型。解决在读数据后显式cast(float)或者在 schema 里直接指定FloatType()。这个错误信息其实很明确但新手容易忽略因为报错发生在 fit 阶段而不是读数据阶段。5.3 rank 设得太大导致单机 OOM现象训练到一半抛OutOfMemoryError: Java heap space或者 Spark 任务卡在某个 stage 不动。原因rank 增大时用户因子矩阵和物品因子矩阵的维度线性增长加上 ALS 每轮迭代要缓存中间结果内存占用是 rank 的倍数。单机 8GB 内存跑 rank200 加 maxIter20很容易撑爆。解决先降 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, truncateFalse) # 定义余弦相似度 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(keylambda x: x[1], reverseTrue) 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()可以保存模型但参数配置要自己写到日志或配置文件里。希望帮到你。本文还有配套的精品资源点击获取
返回列表