资讯详情

基于Spark的电影推荐系统实战:从爬虫到离线召回全链路

📅 2026/10/9 13:43:59 | 华诺云谱 👁 阅读
基于Spark的电影推荐系统实战:从爬虫到离线召回全链路
简介这份资源是一套基于Spark的电影推荐系统完整项目包面向计算机相关专业的在校学生、教师及企业开发者可用于毕业设计、课程设计、项目立项演示或技术进阶学习。项目整合了爬虫数据采集、Web网站、后台管理系统与Spark推荐算法四大模块并附有详细文档形成从数据获取到推荐展示的闭环方案。压缩包共1417个文件约59.61MB涵盖html、css、js等前端页面资源java、scala、py等后端与算法源码以及xml、json、properties等配置文件和parquet数据文件结构完整、层次清晰。目前已有60人学习下载。该资源为个人高分项目源码经导师指导认可答辩评审达95分代码均测试运行成功。读者可据此掌握爬虫实现、前后端交互、后台管理及Spark推荐流程并可在现有基础上修改扩展快速应用于毕设、课设或作业场景。1. 从零搭一套电影推荐系统为什么我劝你先跑通离线召回再碰实时很多人第一次接触推荐系统都是从“协同过滤”四个字开始的然后一头扎进矩阵分解的公式里最后卡在“我到底该用 Spark 的哪个 API”上。这套基于 Spark 的电影推荐系统本质上是一条完整的工业级链路爬虫负责把电影元数据和用户行为抓下来Web 网站承载用户交互后台管理系统做数据维护和推荐结果干预Spark 离线计算引擎负责训练推荐模型并产出召回结果。它解决的核心问题是当你手里只有原始评分和电影信息时如何一步步把它变成一个能上线、能调参、能排查问题的推荐服务。适合谁看如果你已经会写 Python 或 Scala懂一点 SQL但对“推荐系统怎么从数据变成接口”没有完整概念这篇就是给你写的。我不会只讲 ALS 的数学推导而是把爬虫字段设计、Spark 特征工程、离线召回评估、Web 层对接这四段串起来让你照着能跑通一个最小闭环。先记住一个反直觉的结论推荐效果不好八成不是模型选错了而是召回阶段的数据泄漏或负采样出了问题。2. 爬虫与数据层电影元数据和用户评分怎么落成推荐系统的燃料推荐系统的上限在数据层就决定了。很多教程直接甩一个 MovieLens 数据集让你跑 ALS结果一上真实场景就翻车——因为真实数据里用户行为稀疏、电影元数据缺失、时间戳格式混乱。这一章讲清楚爬虫该抓什么、怎么存、怎么清洗成 Spark 能吃的格式。2.1 爬虫字段设计哪些字段真正影响推荐质量我一般把爬虫抓取的字段分成三类。第一类是用户行为字段用户 ID、电影 ID、评分值、评分时间戳。这四个字段是协同过滤的命根子缺一个都会让 ALS 的训练矩阵对不上。第二类是电影内容字段电影 ID、标题、导演、演员列表、类型标签、上映年份、时长。这些字段在冷启动和内容召回里用得上尤其是当某个电影评分记录少于 5 条时纯协同过滤会直接把它忽略这时候类型标签就能兜底。第三类是辅助字段海报 URL、简介文本、语言、地区。这些不直接进模型但后台管理系统和 Web 展示需要。常见做法是爬虫按“电影详情页 用户评分页”分开抓电影详情页用电影 ID 做唯一键评分页用“用户 ID 电影 ID”做联合唯一键。注意时间戳一定要统一成 Unix 秒级或毫秒级整数不要存成字符串否则 Spark 读进来还要额外做类型转换容易在 join 时因为类型不一致丢数据。2.2 用 PySpark 把原始 JSON 清洗成评分矩阵假设爬虫落盘的是每行一个 JSON 的文件字段名可能是user_id、movie_id、rating、timestamp。下面这段代码是我常用的清洗骨架直接跑在spark-shell或 PySpark 脚本里都行。from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, to_timestamp, unix_timestamp from pyspark.sql.types import IntegerType, FloatType, LongType spark SparkSession.builder \ .appName(MovieRecClean) \ .config(spark.sql.shuffle.partitions, 200) \ .getOrCreate() # 读取原始 JSON每行一个评分记录 raw spark.read.json(hdfs:///data/raw/ratings/*.json) # 字段重命名与类型强转防止爬虫字段名漂移 cleaned raw.select( col(user_id).cast(IntegerType()).alias(user_id), col(movie_id).cast(IntegerType()).alias(movie_id), col(rating).cast(FloatType()).alias(rating), col(timestamp).cast(LongType()).alias(ts) ).filter( col(user_id).isNotNull() col(movie_id).isNotNull() col(rating).between(0.5, 5.0) col(ts).isNotNull() ) # 去重同一用户对同一电影只保留最新一条评分 from pyspark.sql.window import Window from pyspark.sql.functions import row_number w Window.partitionBy(user_id, movie_id).orderBy(col(ts).desc()) dedup cleaned.withColumn(rn, row_number().over(w)) \ .filter(col(rn) 1).drop(rn) # 写出 Parquet按评分时间分区方便后续增量训练 dedup.write.mode(overwrite).partitionBy(ts).parquet(hdfs:///data/clean/ratings)这段代码的逻辑说明spark.sql.shuffle.partitions设成 200 是经验值小集群可以降到 50大集群可以升到 500主要影响 join 和窗口函数的并行度。rating过滤在 0.5 到 5.0 之间是因为常见评分体系是半星制如果你爬到的数据是 1 到 10 分制这里要改成between(1, 10)。去重那一步用窗口函数按时间倒序取第一条是为了防止爬虫重复抓取导致同一用户对同一电影有多条评分ALS 训练时如果不去重隐式反馈会被重复计数推荐结果会偏向那些被反复抓到的电影。参数怎么改如果数据量在千万级以下shuffle.partitions设 100 就够如果超过一亿条建议设 500 到 1000并且开启spark.sql.adaptive.enabledtrue让 Spark 自动合并小分区。失败时看什么如果报AnalysisException: cannot resolve ts说明爬虫字段名不是ts去raw.printSchema()里确认实际字段名如果写出 Parquet 时卡在最后一个 task通常是数据倾斜某个热门电影被大量用户评分可以先用salting打散。2.3 电影元数据表与评分表的关联坑电影元数据表通常以movie_id为主键但爬虫抓下来的标题可能带有多语言版本比如同一部电影有中文名和英文名。我一般会保留一个title字段存主标题另外用title_alt存别名在 Web 展示时优先显示主标题。关联评分表时用left join而不是inner join因为有些电影可能还没有评分记录但后台管理系统需要展示它们。如果用了inner join这些电影会直接消失导致后台列表不全。注意电影 ID 在不同数据源里可能一个是整数一个是字符串join 之前务必统一类型否则 Spark 不会报错但结果会是空集这种“静默丢数据”是最难排查的坑之一。3. Spark ALS 推荐模型从训练到离线召回评估的完整链路数据清洗完之后核心就是 Spark MLlib 里的 ALS。很多人调 ALS 只调rank和regParam但真正影响线上效果的是隐式反馈的置信度公式和负采样策略。这一章把训练、评估、召回产出串起来。3.1 ALS 显式反馈与隐式反馈的选型理由显式反馈就是用户直接给的评分1 到 5 星。隐式反馈是用户的行为比如点击、观看时长、收藏。这套系统里爬虫抓到的是评分所以默认走显式反馈。但如果你只有“用户看过这部电影”而没有评分那就得用隐式反馈把implicitPrefs设成True。显式反馈的 ALS 目标函数是最小化评分预测误差隐式反馈则是把行为强度当成置信度行为越多越可信。我一般会先跑一版显式反馈因为评分数据质量高评估指标 RMSE 直观。如果评分数据太稀疏比如只有 1% 的用户有评分那就切到隐式反馈用观看次数或停留时长做置信度。注意隐式反馈下 RMSE 没有意义要看召回率或 MAP。3.2 训练 ALS 模型并调参的代码骨架下面这段是 PySpark 训练 ALS 的完整流程包含训练集测试集切分、参数网格搜索和最优模型保存。from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator from pyspark.sql import functions as F # 读取清洗后的评分数据 ratings spark.read.parquet(hdfs:///data/clean/ratings) \ .select(user_id, movie_id, rating) # 按时间切分用 ts 排序后前 80% 做训练后 20% 做测试 w Window.orderBy(ts) ratings_with_rank ratings.withColumn(rn, F.row_number().over(w)) total ratings_with_rank.count() train ratings_with_rank.filter(F.col(rn) total * 0.8).drop(rn) test ratings_with_rank.filter(F.col(rn) total * 0.8).drop(rn) # 参数网格rank 和 regParam 是必调项 ranks [10, 20, 50] reg_params [0.01, 0.1, 1.0] best_rmse float(inf) best_model None for r in ranks: for reg in reg_params: als ALS( userColuser_id, itemColmovie_id, ratingColrating, rankr, regParamreg, maxIter15, nonnegativeTrue, coldStartStrategydrop ) model als.fit(train) preds model.transform(test) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(preds) print(frank{r}, regParam{reg}, RMSE{rmse:.4f}) if rmse best_rmse: best_rmse rmse best_model model # 保存最优模型 best_model.save(hdfs:///models/als_best)逻辑说明rank是隐向量的维度10 到 50 是常见范围太小欠拟合太大容易过拟合且训练慢。regParam是正则化系数0.01 到 1.0 之间调数据量大时用大一点的正则。maxIter15是迭代次数一般 10 到 20 就收敛了再多收益很小。nonnegativeTrue强制隐向量非负推荐结果可解释性更好但会牺牲一点精度。coldStartStrategydrop是必须的否则测试集里出现训练集没见过的用户或电影时预测值会是 NaNRMSE 直接变 NaN。参数怎么改如果 RMSE 在 0.8 左右徘徊降不下去先检查数据里有没有大量重复评分再去重如果 RMSE 低于 0.5可能是数据泄漏检查测试集是不是混进了训练集。失败时看什么如果报java.lang.IllegalArgumentException: requirement failed: Column user_id must be of type Integer说明 user_id 不是整数回去清洗层强转。3.3 离线召回给每个用户生成 TopN 推荐列表训练完模型只是第一步线上要的是“给用户推荐哪些电影”。ALS 自带recommendForAllUsers但直接用它有个坑它会把用户已经看过的电影也推荐出来。我一般会先拿到推荐结果再和用户历史行为做差集。# 给所有用户推荐 50 部电影 user_recs best_model.recommendForAllUsers(50) # 展开成 (user_id, movie_id, score) 三列 from pyspark.sql.functions import explode flat_recs user_recs.select( user_id, explode(recommendations).alias(rec) ).select( user_id, col(rec.movie_id).alias(movie_id), col(rec.rating).alias(score) ) # 去掉用户已经评分过的电影 seen ratings.select(user_id, movie_id).distinct() final_recs flat_recs.join(seen, on[user_id, movie_id], howleft_anti) # 写出召回结果供 Web 层读取 final_recs.write.mode(overwrite).parquet(hdfs:///data/recall/user_topn)逻辑说明recommendForAllUsers(50)里的 50 是每个用户召回数量线上一般召回 100 到 500 条再交给排序层精排。left_antijoin 是 Spark 里做差集的标准写法比subtract更高效。写出 Parquet 后Web 层可以用 JDBC 或直接读 HDFS 挂载点来获取推荐结果。参数怎么改召回数量根据下游排序能力定如果排序模型简单召回 50 条就够如果排序模型复杂召回 500 条再精排。失败时看什么如果recommendForAllUsers跑得特别慢通常是用户数量太大可以分批跑比如按 user_id 哈希分成 10 批。4. Web 网站与后台管理系统推荐结果怎么变成可交互的产品模型产出的是离线文件用户看到的是网页。这一章讲 Web 层怎么读推荐结果、后台管理系统怎么做推荐干预和数据维护。4.1 Web 层对接推荐结果的三种方式第一种是直连 HDFSWeb 服务挂载 HDFS 目录直接读 Parquet 文件。这种方式简单但延迟高适合离线推荐结果更新频率低的场景。第二种是导入 MySQL 或 PostgreSQL用定时任务把 Parquet 转成表Web 层查 SQL。这种方式查询快但需要额外维护同步任务。第三种是走 Redis 缓存把每个用户的 TopN 推荐结果序列化成 JSON 存进 RedisWeb 层直接GET。我一般用第三种因为推荐结果读多写少Redis 的吞吐量足够而且可以设 TTL 自动过期。常见做法是每天凌晨跑完 Spark 任务后用一个 Python 脚本把final_recs按 user_id 分组写入 Redis 的rec:user:{user_id}键值是电影 ID 列表的 JSON。Web 层拿到列表后再去电影元数据表里查标题和海报。4.2 后台管理系统的推荐干预功能设计后台管理系统不只是增删改查它还要能干预推荐结果。我一般会做三个功能第一是“黑名单”运营可以把某些电影从推荐池里剔除比如版权到期的电影。第二是“加权”运营可以给某些电影手动加权重让它们在召回阶段排更前。第三是“替换”当某个推荐位空着时用运营指定的电影填充。实现上黑名单和加权可以做成两张 MySQL 表Spark 任务在产出召回结果后再读这两张表做过滤和排序。替换逻辑放在 Web 层当 Redis 里取不到推荐结果时走兜底策略。注意后台管理系统的操作要记审计日志谁在什么时候改了哪个电影的权重都要能追溯否则出了问题没法回滚。提示后台管理系统的推荐干预表不要直接改 Spark 的输入数据而是作为后处理层叠加这样模型训练和运营干预解耦排查问题时能快速定位是模型问题还是运营配置问题。5. 避坑与排查这套推荐系统最容易翻车的五个地方这一章是我踩过的血泪经验每条按“现象 → 原因 → 解决”写你照着排查能省很多时间。5.1 现象ALS 训练报 NaNRMSE 评估直接失败原因测试集里出现了训练集没有的用户或电影ALS 默认coldStartStrategy是nan预测值就是 NaN。解决在 ALS 初始化时加coldStartStrategydrop把冷启动样本从评估里丢掉。如果线上必须处理冷启动那就单独走内容召回用电影类型标签做相似度匹配。5.2 现象推荐结果全是热门电影长尾电影一个都不出原因显式反馈下热门电影评分多ALS 会偏向它们如果用了隐式反馈但没做负采样热门电影置信度天然高。解决在召回后做重排序用score / log(1 电影热度)做惩罚或者训练时对热门电影降采样。我一般会在 Spark 里先统计每个电影的评分次数然后在召回结果里过滤掉评分次数超过阈值且分数不高的电影。5.3 现象Web 层读 Redis 推荐结果用户看到的电影和后台管理系统里删掉的电影还在原因Redis 缓存没更新Spark 任务跑完后没有清缓存。解决在 Spark 任务写完 Redis 后给所有rec:user:*键设一个短 TTL比如 25 小时确保每天至少刷新一次。或者用发布订阅机制Spark 任务完成后发一条消息Web 层收到后主动清缓存。5.4 现象爬虫抓下来的评分时间戳是字符串Spark 读进来后排序乱序原因字符串排序是字典序100会排在99前面。解决在清洗层用to_timestamp或cast(LongType)转成数值型再排序。如果时间戳格式不统一比如有的带时区有的不带统一转成 UTC 秒级整数。5.5 现象后台管理系统改了电影权重但推荐结果没变化原因权重表没有被 Spark 任务读取或者读取了但没生效。解决检查 Spark 任务里有没有 join 权重表join 之后有没有按权重排序。另外注意权重表的更新时间和 Spark 任务的调度时间如果 Spark 任务在权重更新之前就跑完了那当然不会生效。我一般会在 Spark 任务启动时打印权重表的最后更新时间方便核对。6. 进阶技巧用离线评估指标反推线上效果最后一章讲一个我常用的验证方法不要只看 RMSE要看召回率和覆盖率。RMSE 低不代表推荐结果好因为 RMSE 衡量的是评分预测精度而推荐系统要的是“把用户可能喜欢的电影排前面”。我一般会算三个指标。第一个是 RecallK在测试集里用户实际评分高于 4 星的电影中有多少比例出现在 TopK 推荐里。第二个是 Coverage所有电影中有多少比例至少被推荐给一个用户。第三个是 Novelty推荐列表的平均流行度越低说明推荐越不依赖热门电影。# 计算 Recall50 K 50 # 测试集里用户高评分电影 test_like test.filter(col(rating) 4.0) \ .groupBy(user_id) \ .agg(F.collect_set(movie_id).alias(liked)) # 推荐结果里每个用户的 TopK recs_topk final_recs.filter(col(score).isNotNull()) \ .groupBy(user_id) \ .agg(F.collect_list(movie_id).alias(rec_list)) \ .withColumn(rec_topk, F.slice(rec_list, 1, K)) # 计算命中率 hit test_like.join(recs_topk, onuser_id) \ .withColumn(hit_count, F.size(F.array_intersect(liked, rec_topk))) \ .withColumn(recall, F.col(hit_count) / F.size(liked)) avg_recall hit.agg(F.avg(recall)).collect()[0][0] print(fRecall{K} {avg_recall:.4f})这段代码的逻辑array_intersect求两个数组的交集size取交集大小再除以用户实际喜欢的电影数量就是单个用户的召回率最后取平均。参数 K 一般设 50 或 100和线上召回数量保持一致。如果 Recall50 低于 0.1说明模型基本没学到东西回去检查数据质量或调大rank。覆盖率计算更简单final_recs.select(movie_id).distinct().count() / 电影总数。如果覆盖率低于 10%说明推荐结果太集中需要加多样性约束比如在召回后做类别打散。我自己的习惯是每次调完 ALS 参数先看 Recall50 有没有提升再看 RMSE。如果 Recall 涨了但 RMSE 没动说明模型在排序上更准了这是好事。如果两个都涨那可能是数据泄漏赶紧查测试集。这套系统从爬虫到推荐接口最耗时的不是写代码而是排查数据不一致和缓存不更新。希望帮到你。本文还有配套的精品资源点击获取
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。

↑