SpringBoot+Spark+Vue电影推荐系统:从日志处理到ALS模型实战
简介一份基于SpringBootSparkVue构建的电影推荐系统完整项目面向计算机科学与技术、软件工程、人工智能等专业学生适用于毕业设计、课程设计、期末作业及答辩演示。项目代码经过导师指导与答辩认可评审得分95分所有功能均测试运行成功下载后可直接部署使用具备一定基础的读者也能基于现有代码二次开发拓展其他推荐场景。压缩包共634个文件包含Java与Scala源码、XML配置、JAR依赖包、Vue前端TypeScript与JavaScript脚本、CSS样式、HTML页面以及少量CSV数据、图片和文档资料整体约66.39MB目录划分清晰便于按后端服务、数据处理、前端展示分层查阅。资源附带了数据流图、系统架构图和详细设计文档能够帮助学习者梳理从数据采集、Spark离线计算、推荐算法实现到前端可视化展示的完整链路。目前已有161人学习下载适合希望快速掌握推荐系统项目实战、完善毕业设计或积累工程经验的人群。1. 电影推荐系统不是堆框架而是把“日志-模型-接口-前端”串成一条线拿到这套基于 SpringBoot Spark Vue 的电影推荐系统源码时我第一反应是看它怎么处理access_log.2019-12-01这份日志。很多毕业设计项目把推荐系统做成“数据库里存了一堆评分然后调一个算法库返回结果”但真正贴近生产环境的做法是从原始访问日志里提取用户行为再用 Spark 做离线统计和模型训练。这个项目好就好在它把access_log到最终 Vue 页面上的“猜你喜欢”完整打通了。适合两类人一类是准备答辩的应届生需要讲清楚推荐链路而不是只讲 CRUD另一类是刚接触 Spark 的 Java 工程师想看看 SpringBoot 怎么和 Spark 任务衔接。下文我会按数据流、ALS 训练、后端接口、前端展示、部署调参这条线拆开讲重点放在可复现的代码和坑位上。2. 数据流设计从 access_log 中提取评分矩阵而不是等用户“打分”2.1 为什么用访问日志而不是直接建评分表常见电影推荐项目会设计一张rating表让用户手动打分。现实中绝大多数用户看完电影根本不会打分但会留下播放、搜索、停留时长等行为。这套资源里的access_log.2019-12-01就是原始行为日志里面包含了用户ID、电影ID、行为类型播放/收藏/评分、时间戳等字段。推荐系统真正要喂给协同过滤算法的是“用户对物品的偏好值”这个值可以从日志里算出来播放算 1 分收藏算 3 分搜索后点击播放算 2 分明确评分则用评分值。这种做法比单独建评分表更接近工业界。如果你直接拿一张现成的评分表做 ALS 训练那只是调库而这个项目让你学会处理原始日志。2.2 日志预处理步骤与代码拿到access_log.2019-12-01后先看格式。常见格式是逗号分隔的六列timestamp,userId,movieId,behavior,duration,rating。我习惯先用 Python 或 Spark 快速探查再决定清洗规则。# 用 Spark 读取并检查日志格式 from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, lit spark SparkSession.builder.appName(MovieLogParser).getOrCreate() df spark.read.option(header, false).option(delimiter, ,) \ .csv(hdfs:///data/access_log.2019-12-01) \ .toDF(timestamp, userId, movieId, behavior, duration, rating) df.show(5, truncateFalse) df.groupBy(behavior).count().show()这段代码用 Spark 读取 HDFS 上的日志文件然后按behavior字段分组统计各类行为数量。里面关键的参数是delimiter ,如果日志是用制表符分隔的就要改成\t。timestamp列用于后续做时间衰减duration可以过滤掉异常数据。提示先用df.show(5)确认字段顺序。很多日志文件前两行是注释需要先过滤。2.3 行为权重映射为评分偏好原始日志里behavior是字符串不能直接用于 ALS。ALS 需要的是userId, movieId, rating三列并且 rating 是浮点数。这里需要写一个权重映射逻辑并考虑时间衰减越近的行为权重越高。# 将行为映射为评分分值并做时间衰减 from pyspark.sql.functions import udf, date_format, current_date, expr, when behavior_score udf(lambda b: {play: 1.0, collect: 3.0, search_click: 2.0, rate_1: 1.0, rate_2: 2.0, rate_3: 3.0, rate_4: 4.0, rate_5: 5.0}.get(b, 0.0), float) df df.withColumn(base_score, behavior_score(col(behavior))) \ .withColumn(days_diff, expr(datediff(current_date(), to_date(timestamp, yyyy-MM-dd)))) \ .withColumn(decay_factor, expr(1.0 / (1.0 0.01 * days_diff))) \ .withColumn(final_rating, col(base_score) * col(decay_factor)) \ .filter(col(final_rating) 0) \ .select(userId, movieId, final_rating)这段代码里behavior_score这个 UDF 把字符串行为映射成分值datediff计算日志日期与当前日期的差值decay_factor让旧行为的分值按指数衰减。final_rating是最终喂给推荐模型的评分。这里有个细节如果你把播放行为算 1 分那么同一个用户看同一部电影十次就相当于打了 10 分会导致推荐结果偏向高频垃圾行为。更合理的是对同一 userIdmovieId 做聚合去重取最大分值或最近一次分值。# 对同一电影的多条行为取最大值避免重复刷分 df df.groupBy(userId, movieId).agg({final_rating: max}) \ .withColumnRenamed(max(final_rating), rating)聚合之后得到训练集。还要注意用户和电影的 ID 在 Spark 中会被解析为字符串或数值如果原始日志里 ID 是 UUID 字符串ALS 要求数值索引。此时需要用StringIndexer将字符串 ID 转为连续整数索引并保存映射关系供后续使用。2.4 数据流图与系统架构的对应关系资源里的系统架构.bmp和数据流图.bmp清晰地画出了完整链路前端 Vue 通过 Nginx 访问 SpringBoot 后端后端接收用户行为日志并写入 MySQLSpark 定期从 MySQL 或 HDFS 读取日志进行离线计算生成推荐结果写回 Redis前端通过/api/recommend接口拿到结果。很多人在答辩时会忽略日志回流这一环前端不只是显示推荐列表还要上报用户行为。在上报接口里需要记录userId、movieId、behavior、timestamp。下面是一个简单的 SpringBoot 日志上报接口。3. 基于 Spark 的 ALS 离线推荐训练、评估与参数选择3.1 为什么选 ALS 而不是 Item-CF 或深度模型ALS交替最小二乘是 Spark MLlib 中实现协同过滤最稳定的算法。相比 Item-CFALS 能处理稀疏矩阵并自动学习隐向量相比深度学习ALS 在中小规模数据集上训练快、可解释性好。这个数据集是单日 access_log量级不大ALS 在 executor 上跑 5 到 10 分钟就能收敛。如果你用自编码器或者 Wide Deep光调特征工程就够写一篇论文了而 ALS 把用户和物品映射到低维向量空间预测评分时做点积。3.2 训练代码与关键参数说明用 Scala 调用 Spark MLlib 是最常见的做法因为 Spark 原生 API 在 Scala 下类型安全度高且不会遇到 PySpark 的序列化问题。这里给出ALSRecommender.scala的核心代码。import org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.SparkSession val spark SparkSession.builder() .appName(MovieALS) .config(spark.executor.memory, 2g) .getOrCreate() val ratings spark.read.parquet(/data/ratings_processed) .select(userId, movieId, rating) val als new ALS() .setMaxIter(15) .setRank(20) .setRegParam(0.05) .setUserCol(userId) .setItemCol(movieId) .setRatingCol(rating) .setColdStartStrategy(drop) val model als.fit(ratings) model.write.save(/models/als_model)这里有几个参数值得细说。rank 20是隐向量维度维度越高模型表达能力越强但也越容易过拟合。对电影推荐这种场景20 到 50 比较常见。maxIter 15是最大迭代次数迭代太少不收敛太多会浪费时间。regParam 0.05是正则化系数调大能让模型更平滑调小会拟合更多噪声。coldStartStrategy drop表示当测试集中出现训练时没见过的用户或电影时直接丢弃预测结果而不是输出 NaN这在计算准确率时特别重要。提示如果训练时报requirement failed: The number of items should be positive说明rating列存在空值或 NaN需要检查预处理阶段是否有用户行为为空的情况。3.3 模型评估RMSE 与覆盖率训练完模型不能直接上线至少要评估 RMSE均方根误差和覆盖率。RMSE 衡量预测评分与真实评分的偏差覆盖率衡量推荐结果覆盖了多少电影。下面的代码把数据按照 8:2 拆分训练测试集。val Array(training, test) ratings.randomSplit(Array(0.8, 0.2)) val model als.fit(training) val predictions model.transform(test) import org.apache.spark.ml.evaluation.RegressionEvaluator val evaluator new RegressionEvaluator() .setMetricName(rmse) .setLabelCol(rating) .setPredictionCol(prediction) val rmse evaluator.evaluate(predictions) println(sRMSE $rmse)RMSE 一般控制在 0.8 到 1.2 之间都是可接受的数值越小越好。如果 RMSE 太大首先检查regParam是否太小导致过拟合其次看rank是否过高。除了 RMSE还要计算推荐列表的覆盖率。Spark 自带RecommendationMetrics不太好用我一般自己算val recommendations model.recommendForAllUsers(10) val totalItems ratings.select(movieId).distinct().count() val recommendedItems recommendations.selectExpr(explode(recommendations) as rec) .select(rec.movieId).distinct().count() val coverage recommendedItems.toDouble / totalItems.toDouble覆盖率达到 30% 以上就算健康。如果你发现覆盖率低可以调大rank或者降低评分阈值让模型给更多物品分配非零向量。3.4 候选物品过滤只推荐未看过的电影ALS 的recommendForAllUsers会推荐用户可能喜欢的 Top N 电影但它不排除用户已看过的电影。如果用户已经看过《肖申克的救赎》再推荐出来就很蠢。过滤逻辑可以在 Spark 做也可以在 SpringBoot 接口层做。推荐在 Spark 里完成这样后端接口就不需要额外查询观看记录。val watched ratings.select(userId, movieId).distinct() val recommended model.recommendForAllUsers(50) val filtered recommended .join(watched, Seq(userId, movieId), left_anti) .select(userId, movieId, prediction)left_anti连接保留左表中右表不存在的记录这样拿到的是用户没看过的电影。注意recommendForAllUsers返回的列名是recommendations需要 explode 后再 join。4. SpringBoot 后端推荐接口设计、参数校验与 Redis 缓存4.1 接口设计用 userId 换取 Top N 推荐后端承担两个职责一是接收前端上传的用户行为日志二是读取 Spark 生成的推荐结果。推荐结果存在 Redis 中键格式为rec:user:{userId}值为电影 ID 列表 JSON 字符串。SpringBoot 接口只需要读 Redis压力很小。RestController RequestMapping(/api/recommend) public class RecommendController { private final StringRedisTemplate redisTemplate; private final MovieService movieService; public RecommendController(StringRedisTemplate redisTemplate, MovieService movieService) { this.redisTemplate redisTemplate; this.movieService movieService; } GetMapping(/{userId}) public ResultListMovieVO recommend(PathVariable Integer userId, RequestParam(defaultValue 10) int limit) { String key rec:user: userId; ListString movieIds redisTemplate.opsForList().range(key, 0, limit - 1); if (movieIds null || movieIds.isEmpty()) { // 兜底策略如果 Redis 中没有推荐结果返回热门电影 return Result.ok(movieService.listHotMovies(limit)); } ListMovieVO movies movieService.listByIds(movieIds); return Result.ok(movies); } }这段代码里RequestParam(defaultValue 10)控制返回条数前端可以按需调整。redisTemplate.opsForList().range从 Redis 列表左侧取 0 到limit-1的元素这样的优点是推荐列表本身是有序的直接按顺序返回即可。Result是统一响应体包含 code、message、data。注意 Redis 里存的 JSON 字符串所以在listByIds时要反序列化。提示不要直接用RequestParam接收 userId 而不做校验负数会导致 Redis 查询异常最终返回 500。4.2 为什么要在 Redis 里放推荐结果Spark 离线计算产生的推荐结果不是实时数据不适合在用户请求时才调用 Spark 计算。一次 ALS 预测要遍历整个模型耗时几百毫秒到几秒这对接口不可接受。正确的做法是把推荐结果预先写入 RedisRedis 读取耗时在 1 毫秒左右。写入逻辑可以在 Spark 训练完成后通过 Spark 的 Redis 连接器完成也可以用 SpringBoot 的定时任务去读取 Spark 输出的 CSV 文件再写入。这个项目里采用的是第二种因为没有把 Spark 和 SpringBoot 进程耦合。定时任务代码如下Component public class RecommendCacheTask { private final StringRedisTemplate redisTemplate; Scheduled(fixedDelay 60 * 60 * 1000) public void cacheRecommendations() { ListRecommendation recommendations recommendationService.loadFromHdfs(); for (Recommendation rec : recommendations) { String key rec:user: rec.getUserId(); redisTemplate.delete(key); redisTemplate.opsForList().rightPushAll(key, rec.getMovieIds()); } } }Scheduled(fixedDelay 60 * 60 * 1000)表示每小时执行一次loadFromHdfs读取 Spark 输出到 HDFS 上的推荐结果目录。rightPushAll会把整批电影 ID 依次推入 Redis 列表顺序保持和 Spark 输出的排序一致。4.3 行为上报接口与日志回写前端每次点击电影详情页都要调用POST /api/behavior把用户行为写入消息队列或直接追加到日志文件。日志文件达到一定大小后由 Spark 自动拉取。PostMapping(/api/behavior) public ResultVoid reportBehavior(RequestBody BehaviorRequest request) { // 参数校验 if (request.getUserId() null || request.getMovieId() null) { return Result.error(400, userId and movieId are required); } // 写入日志文件路径可配置 FileLogWriter.append(String.format(%d,%d,%s,%d, request.getUserId(), request.getMovieId(), request.getBehavior(), System.currentTimeMillis())); return Result.ok(null); }这里有几个关键点行为类型字段使用枚举字符串比如 play、collect、rate_5。如果使用中文如“播放”“收藏”Spark 的when条件要写的分支更多而且编码容易出问题。另外不要直接同步插入 MySQL因为用户点击行为非常高频同步写库会把数据库拖垮。一般做法是先写日志由 Flume 或 Spark Streaming 定期抽取。4.4 后端启动参数与配置application.yml里需要配置 Redis 连接、HDFS 地址、日志文件路径。Spark 的配置放在单独的文件中因为 Spark 任务是独立提交的不依赖 SpringBoot 容器。这里给出常见的配置片段spring: redis: host: 127.0.0.1 port: 6379 timeout: 3000ms lettuce: pool: max-active: 20 max-idle: 10 recommend: hdfs-path: hdfs://localhost:9000/recommend/output log-file: /var/log/movie/access.logrecommend.hdfs-path是 Spark 结果输出目录recommend.log-file是行为日志写入位置。如果部署在本地测试环境可以把 HDFS 路径改为本地目录但 Spark 读取本地文件时要注意每台 executor 节点都要有相同路径否则会报No such file or directory。本地测试时我一般用file:///前缀生产环境用 HDFS。5. Vue 前端推荐列表渲染、M3U8 播放与构建优化5.1 组件结构与样式资源资源里的thumbnail.component.css、star.component.css、mdetail.component.css等文件表明了前端的组件划分。thumbnail是电影缩略图卡片star是评分星星组件mdetail是电影详情页。目录里还有fonts.css、demo.css、app.component.css等全局样式。前端核心逻辑就是一个推荐列表页面调用后端/api/recommend/{userId}接口渲染电影卡片。5.2 推荐列表接口调用与状态管理使用 Vue 2 Vuex 管理用户状态和推荐数据。下面是一个常见的Home.vue中请求推荐接口的代码。template div classrecommend-grid movie-card v-formovie in movies :keymovie.id :moviemovie clickgoDetail(movie.id) / /div /template script import { fetchRecommendations } from /api/movie import MovieCard from /components/Thumbnail export default { name: Home, components: { MovieCard }, data() { return { movies: [], loading: false } }, mounted() { this.loadRecommendations() }, methods: { async loadRecommendations() { this.loading true try { const userId this.$store.state.userId const { data } await fetchRecommendations(userId, 10) this.movies data } catch (e) { // 失败时展示热门电影兜底 const fallback await fetchHotMovies(10) this.movies fallback } finally { this.loading false } }, goDetail(id) { this.$router.push({ path: /movie/${id} }) } } } /script这段代码里fetchRecommendations(userId, 10)是 axios 封装实际请求/api/recommend/{userId}?limit10。接口失败时的兜底策略是请求热门电影接口避免页面白屏。5.3 电影详情页播放 M3U8 流系统里存了电影播放地址很多资源是 M3U8 格式。Vue 播放 M3U8 需要引入video.js或hls.js。热门搜索词里有“vue播放m3u8”这里给出一个简洁可用的播放器组件写法。template div classplayer-container video idvideo-player classvideo-js controls playsinline/video /div /template script import videojs from video.js import video.js/dist/video-js.css import videojs-contrib-hls export default { name: M3U8Player, props: { src: { type: String, required: true } }, mounted() { this.player videojs(video-player, { autoplay: true, controls: true, sources: [{ src: this.src, type: application/x-mpegURL }] }) this.player.play() }, beforeDestroy() { if (this.player) { this.player.dispose() } } } /scriptvideojs-contrib-hls已被废弃新项目建议使用hls.js直接处理但很多老项目仍沿用这套写法。如果你用 Vite 打包videojs-contrib-hls可能会报global is not defined需要在index.html中引入兼容补丁或改用 hls.js。5.4 前端打包与后端部署的注意点Vue 项目打包后是纯静态文件需要放到 Nginx 中。vue.config.js里要设置publicPath: ./否则部署到子路径时资源会 404。推荐接口的代理也要配置避免开发环境跨域module.exports { publicPath: process.env.NODE_ENV production ? / : /, devServer: { proxy: { /api: { target: http://localhost:8080, changeOrigin: true } } } }如果后端接口返回的 JSON 字段是下划线风格如user_id而前端用驼峰可以用 axios 拦截器统一转换或在后端配置 Jackson 的PropertyNamingStrategy.SNAKE_CASE。推荐在后端统一配置前端不用关心字段名映射。5.5 前端样式优先级问题资源里有多个component.css文件Vue 的 scoped 样式不会自动覆盖组件库的样式。你会发现thumbnail.component.css里写的卡片尺寸没生效因为被全局common.css中的同名类覆盖了。解决方法是给每个组件的最外层 DOM 加一个独立的class前缀或者使用深度选择器/deep/ .movie-card调整。注意不同 Vue 版本语法不同Vue 3 中使用:deep(.movie-card)。6. 部署与调参Spark 提交参数、冷启动策略与常见报错6.1 Spark 任务提交命令与资源分配Spark 任务不能直接在 IDE 里跑完就完事推荐用spark-submit提交到 standalone 或 yarn 集群。这里给出一个经过多次调整的提交脚本spark-submit \ --class com.example.recommend.ALSRecommender \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --driver-memory 2g \ --num-executors 4 \ --executor-cores 2 \ --jars mysql-connector-java-8.0.30.jar \ /opt/movie-recommend/als-job.jar \ --input hdfs:///data/access_log.2019-12-01 \ --output hdfs:///recommend/result--executor-memory 4g和--num-executors 4是核心参数。如果你的数据集只有几万条日志给 2 个 executor 就够多了反而慢。--jars用来引入 MySQL 驱动如果 Spark 要从 MySQL 读取训练数据这行必不可少。常见错误是executor-memory设置超过节点物理内存导致 Spark 进程被 YARN 杀掉。建议先用free -h查看可用内存再设置参数。6.2 冷启动处理策略系统上线后一定会有新用户和新电影。用户没有任何行为日志时ALS 模型预测不了recommendForAllUsers不会给这类用户生成结果。这个项目里冷启动策略是返回全局热门电影。热门电影列表可以在 Spark 训练时统计按评分次数和平均评分排序存入 Redis 中固定键rec:hot。SpringBoot 接口在 Redis 查不到用户推荐时就会退回到这个键。另一种常见的冷启动做法是统计近 7 天播放量最高的 20 部电影写一个HotRankJobval hotMovies ratings.groupBy(movieId) .agg(avg(rating).alias(avg_rating), count(rating).alias(cnt)) .filter(cnt 10) .orderBy(desc(avg_rating), desc(cnt)) .limit(20)这个逻辑比单纯按播放量排序更合理播放量高但平均分很低的电影不会冲上来。6.3 日志日期边界问题access_log.2019-12-01只有一天的日志但 Spark 的时间衰减算法依赖datediff(current_date(), timestamp)。如果你现在跑这套代码current_date()是 2026 年计算出的days_diff有两千多天所有行为的分值都被压到接近 0模型几乎训练不出来。这是很多做完项目后发现“推荐结果全是同一部电影”的根因。解决办法是把时间基准固定到日志日期而不是动态取当前日期。可以在预处理阶段把current_date()替换为to_date(lit(2019-12-01))或者直接去掉时间衰减只保留行为权重。项目里如果时间衰减参数写死为 0.01那么在跨年测试时结果退化会非常明显。建议在配置文件里加一个rec.decay.baseDate参数默认值设为日志文件的第一天。6.4 验证推荐结果是否正常模型训练完不要直接看精度指标先手动抽样几个用户看推荐列表是否符合常识。可以写一个简单的验证脚本# 从 Redis 中读取某个用户的推荐列表 redis-cli LRANGE rec:user:1001 0 9如果返回的 10 个 movieId 中有 8 个是同一个系列的电影比如全是漫威说明rank可能太小模型只学到了粗粒度的类型偏好而没有区分具体电影。调大rank到 40重新训练再看。如果返回结果全是同一个年份的冷门电影说明regParam太小模型过拟合到了少数热门样本。另外强烈建议在前端加一个“推荐理由”的弱提示比如“因为你看过《星际穿越》所以推荐《火星救援》”。这不需要复杂的推荐解释算法直接在前端把用户最近观看的电影名称和推荐结果拼一起展示。这样答辩时评委能一眼看出你的推荐不是乱序的随机列表。6.5 部署顺序从零部署这套项目我一般按照这样一个顺序先启动 MySQL 和 Redis再启动 Spark 任务确认模型文件和推荐结果写入 HDFS/Redis最后启动 SpringBoot 后端前端打包到 Nginx。Spark 任务如果失败后端接口依然能返回热门电影兜底不会影响整体演示。这也是这个项目的一个亮点Spark 和 SpringBoot 是松耦合的Spark 挂了不影响 Web 端正常显示只是推荐列表变成热门列表。调试时可以先在本地用 IDEA 跑通 SpringBoot再用spark-submit提交到集群这样问题隔离会比较清楚。本文还有配套的精品资源点击获取