☰
基于Spark的电影推荐系统:ALS协同过滤与全栈链路实战解析
2026/9/26 9:00:23 网站建设 项目流程

简介:基于Spark的电影推荐系统完整工程,整合爬虫数据采集、Web网站展示、后台管理系统及推荐算法核心,面向计算机、人工智能、通信工程等专业在校生与开发者,尤其适合作为毕业设计、课程设计或项目初期演示蓝本。压缩包共含1417个文件,大小约59.61MB,按功能模块清晰组织:前端以html/css/js构建交互页面,后端以Java与Scala实现业务逻辑和Spark计算,Python爬虫脚本负责数据抓取,配合SQL初始化脚本与parquet数据文件完成存储,同时附带详细文档、配置说明及Scala源码,便于快速理解整体架构并搭建运行环境。目前已有59人学习下载。项目代码均经测试运行成功,可完整呈现从数据采集、预处理、推荐计算到结果展示的全流程;除电影推荐主功能外,还提供环境搭建思路、算法实现细节与部署参考,适合具备一定基础的读者直接用于课设、毕设,或在此基础上扩展个性化推荐、实时统计等功能模块。

1. 基于Spark的电影推荐系统:为什么说它是一份能跑的毕设级全栈源码

做推荐系统最容易被忽悠的一点是:你以为难点在算法,实际上难点在数据管道。电影推荐这个场景尤其典型——你既需要真实的用户行为数据,又要把爬虫、存储、推荐算法、Web展示、后台管理串成一条能闭环的链路。这份基于Spark的电影推荐系统,正好把这些环节全部打包了,这也是我拆完第一轮觉得它值得写一篇复现笔记的原因。

资源里包含的不只是Spark推荐代码,而是完整四件套:Python爬虫项目负责采数据、Web网站做前台展示、后台管理系统管内容与用户、Spark推荐系统做ALS协同过滤,外加一份详细文档。对做毕设、课设或者想跑通一个真实推荐链路的人来说,这套代码的参考价值在于它把“数据从哪来、特征怎么算、推荐怎么出、结果怎么展示”整个流程串起来了,而不是只给你一个孤零零的算法文件。

适合谁?计算机相关专业做毕业设计或课程设计的学生,以及想从零跑通一个推荐系统全栈项目、但目前只熟悉其中某一环(比如只会写爬虫、或者只写过Spark单机Demo)的从业者。接下来我会按资源拆解实际可用的技术点,先从系统整体架构和数据链路说起。

2. 系统架构与数据链路:爬虫、Web、后台与Spark各司其职

2.1 四个模块怎么分工:一份数据如何走完推荐全流程

打开这份资源,你能看到四个相互独立又通过数据库串起来的工程。第一块是爬虫项目,用Python抓取电影的基础信息和用户评分记录,这是整个系统的数据入口;第二块是Spark推荐系统,负责读数据、训练ALS模型、产出每个用户的TopN推荐列表;第三块是Web网站,给终端用户提供浏览电影、查看推荐结果的界面;第四块是后台管理系统,让管理员可以维护电影数据、查看用户和推荐状态。

数据流向大概是这样的:爬虫把抓到的电影信息和用户评分写入MySQL,Spark从MySQL或HDFS读取训练数据,训练完成后把推荐结果写回数据库,Web后端查询推荐结果渲染到页面,后台管理系统则直接操作同一份数据库做内容管理。这种设计的好处是模块之间通过数据仓库解耦,任何一个环节要替换实现都比较方便——比如你觉得爬虫质量不够,可以单独重写爬虫而不动Spark部分。

在这类毕设项目的实际答辩中,这个“数据链路闭环”往往是得分的关键。很多学生的项目只有算法没有数据来源,或者只有爬虫没有推荐,而这份资源把业务链路做完整了。MyBatis或MyBatis-Plus注解、Vue或Thymeleaf模板、Spark ALS训练参数,这些都是具体的加分点。

2.2 Spark在这套系统里的位置:离线计算的选型理由

选Spark而不是直接上TensorFlow或PyTorch做推荐,是有明确技术逻辑的。首先,电影评分这种显式反馈数据,最适合的基线算法就是协同过滤里的ALS(交替最小二乘),而Spark MLlib对ALS做了分布式实现,训练数据量大时能横向扩展;其次,Spark天然支持从HDFS、MySQL、本地文件多种数据源读取,与爬虫写入的MySQL数据能无缝对接。

从工程复杂度看,Spark的部署也相对可控。这份资源里的推荐模块核心逻辑是:加载评分数据 → 切分训练集测试集 → 训练ALS模型 → 用模型为每个用户生成推荐列表 → 把结果写回数据库。整个流程用Scala或Java写Spark作业都能实现,本地IDEA里配好Spark依赖就能跑通小数据集,集群部署则是量级上来之后的事。

我一般会建议先把这套代码跑在本地模式(local[*])下,确认推荐结果合理后再考虑提交到集群。本地跑的好处是方便调试和打断点看数据,Spark的local模式已经能模拟分布式执行,对毕设场景来说验证算法逻辑完全够用。等到需要处理百万级评分数据时,再按照资源的集群配置文档去搭Standalone或YARN模式。

2.3 数据模型设计:MySQL表结构怎么支撑推荐与展示

一个容易忽略但很重要的部分是数据表的设计。电影推荐系统的核心表至少有这几张:电影信息表(movie_id、title、genres、release_year等)、用户表、评分表(user_id、movie_id、rating、timestamp),以及推荐结果表。评分表是训练数据的来源,电影表是推荐结果要关联展示的内容,两张表通过movie_id关联。

这里有一个数据建模上的坑:爬虫抓到的数据往往有重复,比如同一部电影被不同页面多次抓取。如果直接入库,训练时ALS会把这些重复记录当成独立评分,导致推荐结果偏向这些重复项。所以入库前需要做去重处理,至少在movie_id这一层做唯一约束。类似的细节资源里的后台管理模块有处理,但你自己复现时最好先检查一遍原始数据质量。

注意:如果资源里自带的数据量很小(几百条评分),ALS训练出的模型可能过拟合。建议先用代码里提供的爬虫脚本补充数据量,再进入训练环节,这样推荐效果才有实际意义。

3. Spark ALS推荐核心:从协同过滤原理到参数调优

3.1 ALS在电影推荐里为什么比普通TopN更合适

ALS(交替最小二乘)是协同过滤中处理显式评分数据最经典的算法之一。它的核心思想是把用户-物品评分矩阵分解成两个低维矩阵——用户特征矩阵和物品特征矩阵,然后用两个矩阵的乘积来预测缺失的评分。所谓“交替”,是因为求解过程中固定一个矩阵优化另一个,交替迭代直到收敛。

在电影推荐场景中,这个思路对应的是“和你口味相似的人喜欢什么,你可能也喜欢”。用户特征矩阵里的每一维可以理解为某种隐式偏好,比如喜欢科幻、喜欢老片、喜欢高评分影片等,这些维度不需要人工标注,算法自己从评分模式中学出来。相比简单的全局热门推荐,ALS能做出个性化结果;相比基于内容的推荐,它不需要提前维护电影标签体系。

Spark MLlib的ALS实现在org.apache.spark.ml.recommendation.ALS包里,直接调用即可。关键参数包括rank(特征维度数)、maxIter(最大迭代次数)、regParam(正则化系数)、alpha(隐式反馈场景的置信度参数,显式反馈时用默认值即可)。这些参数直接决定推荐质量的上下限。

3.2 训练一个ALS模型需要几步:代码级拆解

在Spark工程里,推荐模块的主类一般会经历这样的步骤:创建SparkSession、读取评分数据、数据预处理、划分训练测试集、训练模型、评估指标、生成推荐结果。下面用核心代码片段说明:

// 创建SparkSession,本地模式跑通为主 val spark = SparkSession.builder() .appName("MovieLensALS") .master("local[*]") // 本地模式,所有核参与计算 .getOrCreate() // 读取MySQL中的评分数据(或者从CSV加载,取决于资源中的数据文件) val ratingDF = spark.read .format("jdbc") .option("url", "jdbc:mysql://localhost:3306/movie_db") .option("dbtable", "ratings") .option("user", "root") .option("password", "your_password") .load() // 只保留需要的列,并确保数据类型正确 val training = ratingDF.select( col("user_id").cast("int").as("userId"), col("movie_id").cast("int").as("movieId"), col("rating").cast("float").as("rating") )

这段代码的逻辑很直接:从MySQL读取评分表,把字段类型转成ALS要求的格式。userId和movieId必须是整型,rating必须是浮点型,这是Spark MLlib ALS接口的硬性要求,类型不对会直接报错。实际项目里还会过滤掉评分记录太少(比如少于3条)的用户,减少冷启动和稀疏数据带来的噪声。

接下来是模型训练和预测:

// 切分训练集和测试集,80/20是常见比例 val Array(train, test) = training.randomSplit(Array(0.8, 0.2), seed = 42) // 定义ALS模型并设置参数 val als = new ALS() .setRank(10) // 特征维度,10维起步 .setMaxIter(10) // 迭代10轮 .setRegParam(0.1) // 正则化系数,防止过拟合 .setUserCol("userId") .setItemCol("movieId") .setRatingCol("rating") val model = als.fit(train) // 预测测试集中的评分 val predictions = model.transform(test)

参数这里值得展开说一下。rank决定特征矩阵的维度,太小(比如2)模型表达能力不够,太大(比如100)在小数据集上容易过拟合且训练变慢,MovieLens这种规模的数据集10到20是常见取值。maxIter看收敛情况,一般10到20轮足够,数据量大时可以适当加大,但要注意训练时间会线性增长。regParam控制正则化强度,0.01到0.1这个区间比较常用,过大的正则化会让模型偏向平均值,失去个性化能力。

3.3 评估推荐质量:RMSE与TopN命中率怎么算

模型训练完不能直接说“效果不错”,要用指标量化。ALS最常用的评估指标是RMSE(均方根误差),公式是:对测试集中每个真实评分和预测评分求差的平方,取平均后开方。RMSE越小说明评分预测越准。还有一个业务意义更大的指标是TopN命中率:对每个用户取预测评分最高的K部电影,看其中有多少部是用户实际看过的。

Spark里计算RMSE不需要额外引入库,直接对预测结果做聚合就行:

// 过滤掉预测为空的记录(冷启动用户没有足够数据训练) val validPredictions = predictions .filter(col("prediction").isNotNull) .select("userId", "movieId", "rating", "prediction") // 计算RMSE val rmse = validPredictions .withColumn("squaredError", pow(col("rating") - col("prediction"), 2)) .agg(avg("squaredError").alias("mse")) .select(sqrt(col("mse")).alias("rmse")) .collect()(0)(0) println(s"Test RMSE = $rmse")

TopN命中率则稍微复杂一点,需要把每个用户的预测结果排序取前K,再和真实评分中的高分区求交集。实际项目里常见做法是对测试集里用户真实评分超过3.5分的电影做匹配统计。这两个指标在答辩时非常有说服力,比单说“训练完了”强很多。

3.4 推荐结果落地:写回MySQL供Web端调用

模型只产出预测值还不够,Web网站需要直接查询推荐列表。所以通常会在Spark作业的最后一步,把每个用户得分最高的N部电影写回数据库的推荐结果表:

// 为每个用户生成Top10推荐 val recommendations = model.recommendForAllUsers(10) // 整理成适合写入MySQL的格式 val outputDF = recommendations .select(col("userId"), explode(col("recommendations")).alias("rec")) .select(col("userId"), col("rec.movieId").alias("movieId"), col("rec.rating").alias("pred_rating")) // 写回MySQL的recommendations表 outputDF.write .mode("overwrite") .format("jdbc") .option("url", "jdbc:mysql://localhost:3306/movie_db") .option("dbtable", "recommendations") .option("user", "root") .option("password", "your_password") .save()

这段代码里有个容易踩的细节:recommendations列是数组类型,需要用explode把数组展开成多行才能正常写库。否则你写进去的是一行一个数组,Web端查出来还得做额外解析。我在拆这个项目时发现资源里的原始代码对这一步有处理,但如果你自己改写,很容易在这翻车。

4. Python爬虫与数据入库:从Requests到MySQL的完整链路

4.1 搞清楚爬什么:电影数据源选取与页面结构分析

这套系统的爬虫模块用Python写,目标是抓取电影信息和用户评分。数据源一般是公开的影视网站(比如豆瓣电影或类似的开放式影评站),抓取内容包括电影名称、导演、主演、分类、上映年份,以及用户ID、评分分值。抓取范围不需要太大,对毕设场景来说几百部电影、几千条评分就够训练出能用的模型。

爬虫的第一步是分析目标页面的HTML结构,定位数据所在的标签位置。用开发者工具查看网络请求时,优先找JSON接口而不是直接解析HTML,JSON接口的数据结构更规整,解析代码也更简洁。如果目标站点没有公开接口,再退回到BeautifulSoup或lxml做HTML解析。

注意:这里只说技术实现路径,不讨论特定站点。你自己选数据源时务必检查网站的robots协议和合规性,个人学习项目建议只用公开数据集(比如MovieLens)配合少量爬虫验证流程。

4.2 抓取与解析:Requests加BeautifulSoup的基础套路

爬虫模块的代码结构通常分成三块:请求函数、解析函数、存储函数。请求函数负责拿到页面内容,解析函数负责从HTML提取目标字段,存储函数负责把结构化数据写进MySQL。下面是一个典型的解析函数片段:

import requests from bs4 import BeautifulSoup def parse_movie_page(html): soup = BeautifulSoup(html, "html.parser") movie_list = [] # 定位电影条目所在的DOM节点 for item in soup.select(".movie-item"): movie_id = item.get("data-id") title = item.select_one(".title").text.strip() genres = item.select_one(".genres").text.strip() rating = item.select_one(".rating").text.strip() movie_list.append({ "movie_id": movie_id, "title": title, "genres": genres, "rating": float(rating) }) return movie_list

逻辑说明:用BeautifulSoup的select方法按CSS选择器定位节点,逐条提取需要的字段,组装成字典列表。这里的关键是选择器必须准确对应目标站点的HTML结构,一旦站点改版,这段解析逻辑大概率要重写——这也是爬虫最需要维护成本的部分。

参数和数据清洗的细节也值得注意。rating字段从文本转float时需要处理“暂无评分”这类异常值,genres字段可能是多个分类的组合字符串,入库前要决定是用逗号分隔存一个大字段,还是拆分到关联表。资源里采用的做法是直接存逗号分隔的字符串,这样后续Spark读取时处理简单,但缺点是分类维度的聚合查询会受影响。

4.3 SQLAlchemy入库:为什么不用裸SQL

爬虫产出的数据需要用稳定的方式写入MySQL。如果每条数据都拼SQL字符串执行,代码混乱且容易被注入。SQLAlchemy这样的ORM框架能把Python字典直接映射成数据库记录,代码可读性和可维护性都好很多。下面是一个使用SQLAlchemy的入库示例:

from sqlalchemy import create_engine, Column, Integer, String, Float from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker Base = declarative_base() # 定义电影表结构 class Movie(Base): __tablename__ = "movies" movie_id = Column(Integer, primary_key=True) title = Column(String(200)) genres = Column(String(100)) rating = Column(Float) # 创建数据库连接 engine = create_engine("mysql+pymysql://root:your_password@localhost/movie_db") Base.metadata.create_all(engine) Session = sessionmaker(bind=engine) session = Session() # 批量写入,减少I/O次数 movies = [Movie(movie_id=m["movie_id"], title=m["title"], genres=m["genres"], rating=m["rating"]) for m in parsed_data] session.add_all(movies) session.commit()

这段代码展示了一件事:SQLAlchemy让数据模型和代码里的类一一对应,字段变更时只需要改类定义。批量写入的add_all加commit操作比逐条insert快一个量级——Python和数据库之间的网络往返是主要的性能瓶颈,批量提交能显著减少往返次数。爬虫抓取速度慢没关系,入库慢才是真正让人等得暴躁的问题。

4.4 增量抓取与去重:让爬虫可重复运行

爬虫写完跑一次不是终点。实际使用中需要反复补充数据,比如隔几天更新一次评分。这就要求爬虫支持增量抓取和去重。增量抓取的常见做法是记录上次抓取的页码或时间戳,只抓取新内容;去重则依赖数据库主键或唯一索引。

以电影表为例,movie_id设置为主键后,重复写入相同ID会触发主键冲突。处理方式有两种:一是写入前先查一遍库,存在则跳过;二是使用MySQL的INSERT IGNORE或ON DUPLICATE KEY UPDATE语句。SQLAlchemy中可以用MySQL方言的on_duplicate_key_update,也可以用更简单的方式——先查询判断再插入:

# 查询已有ID集合,只插入新数据 existing_ids = {movie_id for (movie_id,) in session.query(Movie.movie_id).all()} new_movies = [Movie(...) for m in parsed_data if m["movie_id"] not in existing_ids] session.add_all(new_movies) session.commit()

这个方案的缺点是数据量大时内存开销明显,但毕设数据规模一般在几千到几万条,完全够用。如果你要抓百万级数据,建议改用批次查询对比,或者直接用INSERT IGNORE。

5. 部署与避坑:从本地跑通到集群模式的常见问题

5.1 本地环境搭建:IDEA、Maven与Spark依赖的版本匹配

所有代码跑通的前提是环境正确。Spark生态的版本兼容问题是最常见的翻车点,具体来说有三个维度要匹配:JDK版本、Scala版本、Spark版本。比如Spark 3.0以上要求JDK 8或11,且不同版本的Spark编译用的Scala版本可能不同(2.12或2.13),你项目里引入的Scala依赖必须和Spark匹配,否则会出现NoSuchMethodError这类玄学报错。

我一般建议用Maven管理依赖,在pom.xml里显式声明Spark相关依赖。一个能用的依赖声明大致长这样:

<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.12</artifactId> <version>3.1.2</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-mllib_2.12</artifactId> <version>3.1.2</version> </dependency>

artifactId里的_2.12是Scala版本后缀,必须和你的Scala环境一致。资源里的代码如果你用更高版本的Spark去编译,大概率会遇到API废弃或行为变化的问题——比如Spark 3.0后很多DataFrame API的行为有调整,代码不一定能直接跑通。所以建议先用资源的原始版本配置,跑通后再考虑升级。

5.2 必踩的坑:Spark作业运行时的三类典型报错

实际运行Spark推荐模块时,有几种报错几乎必然会碰到一次。我把现象、原因和解决办法整理成一张表,方便你排查时直接对照:

现象原因解决
java.lang.ClassNotFoundException找不到MySQL驱动MySQL JDBC驱动没有加入Spark的classpath把mysql-connector-java依赖加入pom.xml,或用--jars参数指定驱动路径
org.apache.spark.SparkException: Task not serializable在RDD或DataFrame的算子中使用了未被序列化的类实例(比如在map里调用了外部对象的方法)把使用到的对象改成static/object,或在算子内只使用局部变量
AnalysisException: Path does not exist或找不到表JDBC连接串格式错误,或MySQL表名/字段名大小写不匹配检查URL格式jdbc:mysql://host:port/dbname,确认表名和执SQL时的名称完全一致

Task not serializable是初学者最常困惑的报错之一。原因是Spark的分布式执行需要把闭包里的引用对象序列化后发送到多个Executor执行,如果你的闭包里引用了一个不可序列化的类——比如含数据库连接的对象——就会报错。解决办法很简单:不要在算子内部创建或引用数据库连接,连接只放在Driver端建立,或者用foreachPartition在分区内部单独创建连接。

5.3 Web与后台联调:推荐结果展示为空或数据不刷新

Spark作业把推荐结果写回MySQL后,Web端去查询这些结果。最容易出现的问题是:推荐表结构设计不兼容,或Web端查询SQL写了错误的字段名。比如Spark写回的字段名是pred_rating和movie_id,Web端的Mapper里对应的是predRating和movieId,如果没做字段映射,查出来结果全为空。

另一个常见问题是数据不刷新。Spark作业重跑后推荐表更新了,但Web页面看到的还是旧数据。这通常是因为Web应用做了缓存,或者浏览器页面缓存了接口响应。排查时先清缓存看接口返回,确认接口返回的是新数据,再排查前端渲染层。如果是Spring Boot应用且开了缓存注解,可以在Mapper方法上暂时去掉缓存注解验证。

还有一个不得不提的坑:时区问题。Spark写时间戳到MySQL时,如果Spark和MySQL的时区不一致,推荐结果表中记录的时间字段会差8小时。虽然不影响推荐结果本身,但在后台管理系统的“最后推荐时间”这个展示字段上会非常显眼。解决办法是在JDBC连接串上加serverTimezone=Asia/Shanghai参数。

5.4 Spark集群搭建的注意事项

资源文档里包含Spark集群搭建的步骤,实际部署时和本地跑有很多不同。首先是内存配置,Spark作业的Executor内存默认1G,如果你的评分数据是百万级,训练时会出现OOM。我一般会把Executor内存调到2G到4G,Driver内存根据本机内存情况设置。无论是Standalone模式还是YARN模式,配置都在spark-env.sh或提交命令的--executor-memory参数里。

其次是文件依赖问题。Spark作业如果依赖外部配置文件或多模块的JAR包,提交集群时需要把依赖打进去或用--files分发,否则作业会在运行时找不到类或配置。血的教训是:本地跑通不代表集群能跑通,提交前先检查依赖打包是否完整。如果资源里的文档没写清楚这一步,按mvn clean package打包后通过spark-submit --files把资源文件一起提交,是稳妥的做法。

6. 推荐效果的验证技巧与冷启动处理

模型训练完推荐结果也写库了,怎么判断这套系统真的有效?很多人看一眼推荐列表觉得“还行”就过了,这样在答辩时很容易被追问卡住。我常用的办法是交叉验证加评分分布分析。交叉验证就是把数据切成K份轮流做训练和测试,取平均RMSE作为最终评价,比一次性划分更稳定。评分分布分析则是看推荐结果中是否有某个头部电影被推给了几乎所有用户——如果出现这种情况,说明模型没有学到个性化,只是隐性全局热门榜。

用Spark做K折交叉验证不需要自己写循环,直接用MLlib里的CrossValidator即可,但要注意ALS的计算复杂度,K值取5比较平衡。除了数值指标,还可以做业务层面的验证:随机选几个用户,打印他评分过的电影和推荐出的前10部电影,看是否有重合、类型是否相关。这一步的观察结果写进答辩PPT里,比贴一个RMSE数字更有说服力。

冷启动处理是这类型项目最容易暴露深度的地方。新用户没有任何评分记录,ALS无法为他生成用户特征向量,推荐结果为空或退化到热门推荐。资源里的系统对这个问题有处理,但如果你要改进,常见的思路是做一个简单的规则推荐:新用户注册后先用全局高分榜或热门榜填充推荐列表,等他产生3到5条评分后再替换成ALS个性化结果。这个策略实现简单,效果立竿见影。

提示:验证推荐效果时,不要只看RMSE绝对值。RMSE受数据稀疏程度影响很大,评分稀疏时RMSE天然偏大。要和随机预测或全局平均值预测的RMSE做对比,才能真正看出ALS模型提升多少。

数据库层面也有一个我反复提醒自己的习惯:每次重跑Spark作业前,先确认推荐结果表是否被之前的作业锁住或存在脏数据。如果上次作业中途失败,可能留下半张表数据,与新的推荐结果混在一起。所以我每次都会在作业脚本开头加一步清空推荐表的操作,确保写库的实时性和数据的一致性。从那以后,每次重跑推荐任务我都强制先执行一次TRUNCATE recommendations再跑作业,宁可多花几秒钟,也不让旧数据干扰结果判断。希望这个习惯能帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询