☰
SpringBoot+Spark旅游推荐系统毕设实战:ALS协同过滤与离线推荐全解析
2026/10/8 3:20:13 网站建设 项目流程

每年到了毕业设计的季节,总有一大批同学纠结选什么题目。推荐系统这个方向热度一直很高,但很多人做着做着就做成了“增删改查管理系统”,推荐逻辑就写个if-else,答辩时被老师一问就露馅了。“SpringBoot + Spark的旅游推荐系统”这个课题,如果你能真正把Spark用起来,把推荐算法落地,那效果完全不一样。这篇文章我把整个项目的设计思路、算法原理、数据库结构、核心代码、踩坑记录全部拆开讲,想选这个题或者正在做的同学可以直接照着搭,不用再自己乱撞了。

1. 项目到底在做什么,为什么要用Spark

1.1 课题定位与核心功能拆解

我理解这个题目想要的东西是这样的:用户登录系统之后,系统能根据他之前的浏览、收藏、搜索、评分等行为,从景点库里找出他大概率感兴趣的景点推荐给他。同时系统要有完整的后台管理、景点信息展示、用户行为记录、推荐结果展示这些基础功能。

说白了,这是一个带有“智能推荐”能力的旅游信息管理平台,而不只是普通的信息管理系统。核心亮点就在推荐模块,这也是能跟普通CRUD拉开差距的地方。整个系统拆开来看大概有这几块:

  • 用户端:注册登录、浏览景点、搜索景点、收藏景点、评分、查看推荐列表
  • 管理端:景点管理、用户管理、标签管理、行为数据统计、推荐策略配置
  • 推荐引擎:基于用户历史的离线推荐 + 冷启动兜底推荐

1.2 为什么方案选型是SpringBoot + Spark

很多同学会问:做一个推荐系统,用SpringBoot不就行了,为什么非要引入Spark?这也正是这个课题最有价值的地方。

先说SpringBoot。它是Java生态里做Web后端最顺手的框架,内置Tomcat、Spring MVC、MyBatis这些,开发效率很高,适合快速把管理端和用户端页面跑起来。毕业设计用它做业务底座完全没有问题。

但推荐系统的核心是“算”。如果用户量小、数据量小,用单机的Python或者简单的SQL聚合就能算完。但如果你要演示的是一个有大数据背景的推荐系统,Spark就有它不可替代的作用了——它是分布式的计算引擎,能把“海量用户行为数据 + 海量景点数据”的矩阵运算、协同过滤训练放到集群里去跑。

而且这个题目用Spark还有一个“讨巧”的点:你可以把推荐计算做成离线批处理任务,定期执行。比如每6小时跑一次推荐任务,训练ALS模型,计算所有用户的TopN推荐列表,把结果写到MySQL或Redis里,前端直接查结果就行。这样既用上了Spark的分布式计算能力,又不至于像纯实时推荐那样把架构搞得特别复杂。这是毕设最稳妥的玩法,也是真实工业界最常见的做法。

1.3 整体架构与数据流转全过程

我画一个整体的数据流转流程,你就能非常清晰地理解系统长什么样:

用户在小程序/网页端产生行为(浏览、点击、收藏、评分)→ 行为数据写入MySQL → Spark定时任务从MySQL拉取行为数据 → 数据清洗、构建用户-景点评分矩阵 → ALS模型训练 → 生成所有用户的TopN推荐结果 → 推荐结果写回Redis和MySQL → SpringBoot查询接口从Redis读取推荐结果 → 前端展示

这个流程里,业务系统和计算系统是解耦的。SpringBoot管业务,Spark管计算,两者通过数据库和中间件通信,不互相干扰。这就是为什么说它是一个“智慧旅游个性化推荐平台”——它的核心计算引擎是独立的大数据组件,不是塞在Controller里的一段循环代码。

2. 推荐算法核心原理,这一块必须吃透

2.1 协同过滤的思路:用“同好”来推荐

推荐系统算法里最经典的就是协同过滤。它的核心思想用一句话概括就是:如果你的行为和某个用户群体很相似,那就把那个群体喜欢的东西推荐给你。

比如你和10个用户都收藏过“故宫”“颐和园”“南锣鼓巷”,那这10个用户还收藏了“天坛”,系统就会把“天坛”推荐给你。这就是基于用户的协同过滤。反过来,如果你看过的景点和“故宫”相似(喜欢故宫的人也喜欢颐和园),那系统就给你推荐“颐和园”,这是基于物品的协同过滤。

在实现层面,这两类都需要计算相似度矩阵。基于用户的相似度计算复杂度是O(N²)级别,用户量一大就扛不住。所以工业界更常用的是基于物品的协同过滤,或者更进阶的矩阵分解方法。在Spark的MLlib里,最常用的就是ALS(Alternating Least Squares,交替最小二乘法),这属于矩阵分解家族的算法。

2.2 ALS矩阵分解到底在干什么

矩阵分解你别被名字吓到,我用一个很通俗的类比来解释。

假设系统里有1000个用户、500个景点。我们把用户对景点的评分想象成一张巨大的表格:行是用户,列是景点,格子是评分。但绝大多数格子是空的,这叫做评分矩阵的稀疏性。ALS干的事情,就是把这个稀疏矩阵拆成两个小矩阵相乘:

  • 用户特征矩阵U(1000行 × K列),每一行代表一个用户被压缩成K个隐式特征
  • 物品特征矩阵V(500行 × K列),每一行代表一个景点被压缩成K个隐式特征

K一般取10到50之间。用户向量和景点向量的点积,就是预测评分。ALS算法通过不断交替固定U优化V、固定V优化U的方式,让预测评分尽量接近真实评分。这个K个特征不需要我们定义具体含义,它可能是“历史人文程度”“自然风光比例”“亲子友好度”这种可解释的,也可能是无法解释的隐含因子。

在Spark里,ALS的直接用法是这样的:

// Scala版 Spark使用ALS训练推荐模型 import org.apache.spark.ml.recommendation.ALS val als = new ALS() .setMaxIter(10) // 迭代次数 .setRegParam(0.01) // 正则化参数,防止过拟合 .setRank(10) // 隐特征数量K .setUserCol("userId") .setItemCol("scenicId") .setRatingCol("rating") val model = als.fit(trainingData)

这段代码看起来简单,但有几个点需要特别注意。setMaxIter太小模型不收敛,太大训练时间翻倍,10到15通常够用。setRegParam如果设置太大,模型会偏向“平均水平”,推荐结果没有个性化;设置太小则容易过拟合训练数据,预测新数据效果差。setRank更是个经验活,拿一个验证集反复试,哪个K使得RMSE(均方根误差)最低,就选哪个。

2.3 评分数据从哪里来:显式反馈与隐式反馈

很多同学把评分理解成“用户打了多少分”,这是对的,但不够全面。实际系统里,评分来源大概有三类:

第一类是显式评分。用户进来给景点打分,1到5分,这个分数最靠谱。但问题是绝大部分用户不会主动打分,数据稀疏性非常严重。

第二类是隐式行为转化。用户收藏了一次算3分,搜索点击算1分,浏览详情超过1分钟算2分,分享算5分。把这些行为按权重累加成一个综合得分,能极大丰富评分矩阵。我在实际项目中用过一个转化权重表,效果很好:

行为类型权重说明
景点浏览1浏览详情页超过30秒
景点收藏3点击收藏按钮
景点分享5分享到社交平台
景点评分用户评分值显式分值直接取用
搜索关键词命中2搜索结果点击景点

第三类是默认值填充。对于完全没有任何行为的用户,用该景点的平均评分做兜底,配合后续规则推荐。

我个人建议,毕设里显式评分和隐式转化两种都要做进去。只做显式评分,你的数据会非常稀疏,Swiss cheese一样的矩阵根本训练不出好模型;只做隐式转化,答辩时老师会问你“凭什么给这些行为定这个权重”,你还需要一套设计依据来回答。

2.4 多路召回:不要只靠一个“黑盒”模型

还有一个非常关键的工程化思路:推荐结果不要只从ALS模型输出。业界一般用“多路召回”的策略——多种算法或规则分别召回一批候选景点,最后合并、去重、过滤、排序,形成最终列表。

ALS协同过滤只是其中一路。你还要加:

  • 基于内容的召回:把景点的标签(历史古迹、自然风光、亲子游、美食打卡等)做成TF-IDF向量,算景点之间的相似度,用户喜欢过某个景点,就把相似景点推荐出来。这种召回结果对“新景点”特别友好,因为它不依赖用户评分。
  • 热门榜兜底:统计全站浏览和收藏最高的前20个景点,作为保底推荐。这个对冷启动用户特别重要。
  • 地域偏好:如果用户所在的省份有多个景点,可以加权重倾斜。

我是怎么做的呢?ALS算出每个用户对500个景点的预测评分,取Top50;内容相似度算出Top50;热门榜取Top20;合并起来去重,然后按预测评分加权排序,最终取Top10存到Redis。这样推荐列表的覆盖率和准确性都有保障,答辩时这个“多路召回”设计本身就是个加分项。

3. 数据库设计与后端模块划分

3.1 核心数据表设计,一次讲透

推荐系统的数据库设计和普通管理系统有个很大的区别:你要多存“行为流水”和“推荐结果”两张关键表。下面我给出最核心的表结构和设计理由。

用户表t_user:主键id、用户名、密码、昵称、性别、省份、城市、注册时间。省份和城市字段非常重要,地域偏好召回要用,很多同学会忽略。

景点表t_scenic:主键id、景点名称、景区等级(5A/4A)、所在省份、所在城市、票价、开放时间、景点描述、封面图URL、标签列表。标签可以用逗号分隔的字符串存,也可以用独立的t_tag和t_scenic_tag关联表,更规范。毕设为了快速开发可以用逗号分隔,但用关联表更经得起追问。

关键行为表t_behavior:主键id、userId、scenicId、行为类型(1浏览、2收藏、3分享、4评分)、行为分数、行为时间。这张表是Spark任务的数据来源,一定要注意查询索引,联合索引(userId, scenicId)必加,否则数据量稍微大一点查询就慢得离谱。

评分表t_rating:主键id、userId、scenicId、评分值、评分时间。这张表可以看作是对行为表的规整汇总,Spark直接读它会更容易构建评分矩阵。

推荐结果表t_recommend:主键id、userId、scenicId、推荐评分、推荐来源(1协同过滤、2内容相似、3热门榜、4混合)、排名、生成批次号、创建时间。批次号字段很重要,每次跑推荐任务时生成一个batchId,方便回滚和对比不同版本模型的效果。

3.2 SpringBoot后端模块怎么划分

后端我建议按功能垂直切分,包结构清晰一点,答辩时有条理:

  • controller:接收前端请求,返回JSON
  • service:业务逻辑层,推荐结果查询、用户行为落地、后台管理逻辑
  • mapper:MyBatis持久层
  • recommend:与Spark交互相关的类,比如调用Spark任务、读推荐结果、缓存管理
  • model:实体类,对应数据库表

举个例子,推荐接口的调用链是:前端请求/recommend/list?userId=1→ Controller接收 → Service层先查Redis缓存,缓存key设计成recommend:user:{userId}→ 如果Redis查到就直接返回;查不到就去t_recommend表按batchId最新一批查;查到后回填Redis并设置过期时间,比如6小时。

这里面的一个核心技巧是:Spark算完之后你要保证结果在Redis里能快速被查询,避免每次推荐请求都去MySQL里做多表关联查询。Redis的value建议直接存JSON字符串列表,取出后反序列化成对象,前端拿到的接口延迟基本在30毫秒以内。

3.3 用户行为记录接口的一个设计细节

用户行为数据是整个推荐系统的“汽油”,这块没做好,后面算法全是空转。我强烈建议你把行为上报接口做成统一入口:

// 统一行为上报接口 @PostMapping("/api/behavior") public Result addBehavior(@RequestBody BehaviorDTO behaviorDTO) { // 1. 校验参数 // 2. 根据行为类型计算分数 // 3. 插入 t_behavior 流水表 // 4. 如果有显式评分,同步写入 t_rating // 5. 返回成功 }

前端在页面上的所有操作——浏览景点详情、收藏、评分、搜索点击——都请求这个接口,把行为类型传过来。这样Spark任务只需要增量读取t_behavior表,不用关心页面上各种复杂的埋点逻辑。增量读取也很简单,t_behavior表加一个自增id,Spark记录每次跑批时的最大id,下次只拉id > 上次最大id的数据,这就是最基础的增量计算方案。

4. Spark离线推荐模块,一步一步说清楚

4.1 环境准备与项目结构

做Spark开发,我强烈建议直接用Scala编写Spark作业,然后打成一个Jar包,通过Java命令或者Spark-submit提交运行。但考虑到毕设项目要和SpringBoot互通,也可以用纯Java编写Spark作业,这样可以在同一个Maven项目里管理。我个人的建议是:如果对Scala不熟悉,就全用Java写,不要给自己挖坑。

Maven依赖大致是这些:

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

这里scope用provided很重要,因为实际运行时Spark集群环境自己会带这些依赖,如果你打入Jar反而会和集群自带的版本冲突。本地调试的时候IDEA会从编译classpath里找到依赖,也没问题。

Spark版本的选择也有讲究。3.x版本和2.x版本API变化比较大,我建议选3.1.2或3.2.1,网上资料最多,踩坑容易找到解决方案。Java版本配合Spark 3.x需要Java 8或Java 11,别用Java 17太新版本,会出现一些莫名其妙的兼容问题。

4.2 从MySQL拉取数据并构建评分矩阵

Spark作业的第一步是读取MySQL里的用户行为数据。这一步如果你直接select * from t_behavior,当表数据量到几十万行时,这个并行度就有问题——Spark默认可能只有一个分区在读数据,速度很慢。你可以手动设置分区策略。

下面这段是核心代码,做了分区读取优化:

// Spark SQL从MySQL读取行为数据,并做分区优化 Map<String, String> options = new HashMap<>(); options.put("url", "jdbc:mysql://127.0.0.1:3306/travel_rec?useSSL=false&serverTimezone=Asia/Shanghai"); options.put("dbtable", "t_behavior"); options.put("user", "root"); options.put("password", "root"); options.put("driver", "com.mysql.cj.jdbc.Driver"); // 按id字段分段并发拉取,提高读取效率 options.put("partitionColumn", "id"); options.put("lowerBound", "0"); options.put("upperBound", "1000000"); options.put("numPartitions", "4"); Dataset<Row> behaviorDF = spark.read() .format("jdbc") .options(options) .load();

这里的partitionColumn要用自增主键,numPartitions设成执行器数量的整数倍,可以让数据均匀分片读取。如果不用这个配置,Spark会单机单线程把整张表拉回来,数据一多直接卡死。

接下来要把行为数据映射成评分值。如果评分表t_rating里有显式评分,直接优先用;否则用行为权重表折算。Spark里可以用SQL表达式一次搞定:

behaviorDF.createOrReplaceTempView("behavior"); Dataset<Row> ratingDF = spark.sql( "SELECT " + " userId, " + " scenicId, " + " MAX(score) AS rating " + "FROM ( " + " SELECT userId, scenicId, " + " CASE behaviorType " + " WHEN 1 THEN 1 " + " WHEN 2 THEN 3 " + " WHEN 3 THEN 5 " + " WHEN 4 THEN behaviorScore " + " ELSE 1 " + " END AS score " + " FROM behavior " + ") temp " + "GROUP BY userId, scenicId" );

MAX是聚合函数,目的是如果用户对同一个景点既有浏览又有收藏,就取权重最高的那个行为对应的分数。这样就得到了一个完整但稀疏的“用户-景点-评分”三元组。

4.3 ALS训练与调参,R值怎么定

ALS在Spark MLlib里的调用,我前面已经放了Scala版。这里放一个完整的Java版,包含训练、预测、评估:

// Java版ALS训练与评估 import org.apache.spark.ml.evaluation.RegressionEvaluator; import org.apache.spark.ml.recommendation.ALS; import org.apache.spark.ml.recommendation.ALSModel; Dataset<Row>[] splits = ratingDF.randomSplit(new double[]{0.8, 0.2}, 42L); Dataset<Row> training = splits[0]; Dataset<Row> test = splits[1]; ALS als = new ALS() .setMaxIter(15) .setRegParam(0.01) .setRank(20) .setUserCol("userId") .setItemCol("scenicId") .setRatingCol("rating"); ALSModel model = als.fit(training); // 预测测试集评分并计算RMSE Dataset<Row> predictions = model.transform(test); RegressionEvaluator evaluator = new RegressionEvaluator() .setMetricName("rmse") .setLabelCol("rating") .setPredictionCol("prediction"); double rmse = evaluator.evaluate(predictions); System.out.println("RMSE = " + rmse);

注意一个容易踩坑的点:ALS预测评分的时候,如果测试集里的用户或景点是训练集中从未出现过的,预测结果里会出现NaN值。这些NaN在计算RMSE时会导致evaluator直接报错,或返回NaN。处理方式是过滤掉预测值为NaN的行:

Dataset<Row> filtered = predictions.filter("prediction IS NOT NULL AND prediction != 'NaN'");

关于rank参数的取值,我的经验是从5到50每隔5取一个值分别训练,记录下来哪个rank使RMSE最低,就选哪个。rank太小,特征表达能力不够,模型欠拟合;rank太大,特征能表达得很细,但容易过拟合,且计算量显著增加,对海量稀疏矩阵来说rank=10~30基本够用。

4.4 生成推荐列表并写回Redis和MySQL

模型训练完之后,不能只停留在预测评分阶段,要计算出“每个用户最可能的TopN个景点”。Spark提供了现成的recommendForAllUsers方法,能直接为所有用户生成TopN推荐:

// 为所有用户生成Top10推荐列表 Dataset<Row> recommendations = model.recommendForAllUsers(10);

这个返回的DataFrame每行是一个userId加一个数组结构,数组里是{scenicId, rating}。你需要把它展开成扁平的行记录,才能方便写入MySQL:

// 将嵌套数组展开成(userId, scenicId, score)格式 Dataset<Row> flattened = recommendations .select( col("userId"), explode(col("recommendations")).alias("rec") ) .select( col("userId"), col("rec.scenicId").alias("scenicId"), col("rec.rating").alias("score") );

explode是Spark里的“炸裂函数”,作用是把一个数组字段拆成多行。如果你第一次见到这个函数,记住这个场景——它就是为了处理这种“一行里包含一个列表”的数据而生的。

扁平化之后,就可以批量写入MySQL了。但这里有个性能问题:如果循环逐条插入,500个用户×10条推荐就是5000条数据,逐条insert可能要几分钟。一定要用批量插入:

// 使用JDBC批处理写入,性能提升明显 flattened.foreachPartition((Iterator<Row> iterator) -> { Connection conn = getJdbcConnection(); // 获取数据库连接 conn.setAutoCommit(false); PreparedStatement ps = conn.prepareStatement( "INSERT INTO t_recommend(userId, scenicId, score, source, batch_id) VALUES (?,?,?,?,?)" ); while (iterator.hasNext()) { Row row = iterator.next(); ps.setLong(1, row.getLong(0)); ps.setLong(2, row.getInt(1)); ps.setDouble(3, row.getDouble(2)); ps.setString(4, "ALS"); ps.setLong(5, batchId); ps.addBatch(); } ps.executeBatch(); conn.commit(); ps.close(); conn.close(); });

foreachPartition的意思是对每个Spark分区,只在分区内建立一个数据库连接,然后把这一个分区内的所有行批量插入。这样连接数少,批量提交快,整体写库效率远高于逐行插入。这个细节在实际生产环境几乎是标配,你写进毕设里也是一个亮点。

Redis侧的写入其实可以和MySQL并行做。按用户维度批量写入:

// 按用户维度写入Redis,key是recommend:user:{userId} flattened.foreachPartition(iterator -> { // 每个分区内按userId聚合后使用pipeline批量写入Redis });

考虑到毕设的教学目的,我会把MySQL作为权威数据源,Redis只做查询加速。如果Redis挂了,系统直接查MySQL兜底,不影响基本功能。

4.5 SpringBoot侧查询推荐结果接口

用户端请求推荐时,SpringBoot本身并不做计算,只做查询。查询策略是:优先Redis缓存,命中失败则查MySQL最新批次。这个策略在上文提到过,对应接口代码大概长这样:

// 推荐查询接口核心逻辑 public List<ScenicVO> getTopNRecommend(Long userId, int limit) { // 1. 先查Redis String key = "recommend:user:" + userId; String cached = redisTemplate.opsForValue().get(key); if (StringUtils.hasText(cached)) { return JSON.parseArray(cached, ScenicVO.class); } // 2. Redis未命中,查MySQL最新批次 Long latestBatchId = recommendMapper.getLatestBatchId(); List<RecommendPO> recs = recommendMapper.selectByUserAndBatch(userId, latestBatchId, limit); // 3. 查景点信息并组装VO List<ScenicVO> result = parseScenicInfo(recs); // 4. 回填Redis缓存 redisTemplate.opsForValue().set(key, JSON.toJSONString(result), 6, TimeUnit.HOURS); return result; }

缓存过期时间设6小时比较好,因为Spark离线任务最频繁也就是每6小时跑一次,缓存比这个周期长没有意义,比这个周期短又浪费数据库查询。

5. 冷启动、实时推荐与效果评估

5.1 新用户和新景点的冷启动处理

ALS模型有一个先天缺陷:对于没有行为记录的“新用户”,它没有办法算出该用户的特征向量,推荐结果自然无从谈起。系统里随便来个注册30秒的用户,如果他从来没有点过任何景点,推荐接口返回空列表,这不就尴尬了。

我用的冷启动方案分三层。

第一层,针对新用户,直接返回热门榜Top20。热门榜怎么定义?我推荐用最近7天所有用户行为流水,按景点聚合统计行为次数,按权重求和倒序排。要写SQL也很简单:

SELECT scenic_id, SUM(score) AS hot_score FROM t_behavior WHERE create_time > DATE_SUB(NOW(), INTERVAL 7 DAY) GROUP BY scenic_id ORDER BY hot_score DESC LIMIT 20;

第二层,当用户有几个浏览行为之后,系统尝试做内容相似推荐。用户看过的景点打上标签,比如“古迹”“山水”“美食”,然后找出同标签下热门的其他景点。

第三层,把“新景点”单独拎出来做一个“为你探索”一栏,随机从最近上架、还没太多人看过的景点里挑几个展示。这个三十年前的概念在业界有据可查,叫探索与利用权衡。你可以在老用户推荐列表里混入少量新景点,前面用ALS模型算的准确率高,但太准了也会让人审美疲劳。

我实际项目里的混入比例是:前8个推荐位放协同过滤结果,第9个和第10个放一个随机新景点和一个地域相近景点。

5.2 实时推荐扩展思路,答辩常问

毕设的离线推荐已经能撑住场面,但如果答辩老师问“用户刚点了一个景点,怎么立刻给他推荐相似的?”你得能答上来。

实时推荐的经典做法是:用户行为发生时,先用Redis记录最近行为列表,再根据景点标签的相似度,从预先计算好的“景点相似度矩阵”里取出TopN。在SpringBoot里,这个逻辑可以直接查一张t_scenic_similarity表,这张表是Spark离线任务算好的,每行存(scenic_id, similar_scenic_id, similarity_score)。

用户点击景点A后,后端查相似表,找出和A相似度最高的几个景点,过滤掉用户已经看过的,直接拼一个实时推荐列表返回。这个表可以用Redis的Set或ZSet存,也可以用MySQL查,量不大都很快。当你讲出这个流程,不仅证明你理解了推荐系统的实时链路,还说明你有工程落地能力。

5.3 效果评估:怎么证明你的推荐系统“有效”

很多毕设做完推荐系统,却说不出它到底哪里好。认真做效果评估,能从一堆论文式空谈里脱颖而出。

最简单直接的指标是离线RMSE,这个训练时已经算过。但RMSE低只说明模型预测用户评分准,不代表用户体验好。建议再统计几个在线指标:

  • 推荐位点击率:前端埋点统计推荐列表里景点的曝光次数和点击次数
  • 推荐位转化率:从推荐列表点进去并最终收藏的比率
  • 个性化离散度:不同用户得到的推荐列表重合度,如果所有用户的推荐结果几乎一样,说明个性化程度不足

我在项目里增加了一张t_recommend_log表,专门记录“哪个用户、在什么时间、推荐了哪些景点、点击了哪个”。每有推荐曝光时就写一条日志,点击时再更新状态。答辩时,你拿出一张统计图,展示推荐位点击率从初版模型的8%提升到优化后的15%,这个说服力比任何口头解释都强。

6. 实战踩坑记录与性能调优笔记

6.1 依赖冲突与版本兼容,最让人头疼

Spark与SpringBoot集成时,最容易崩的地方是依赖冲突。SpringBoot 2.x依赖的jackson版本和Spark 3.x依赖的jackson版本不一致,运行时就报NoSuchMethodError或者ClassNotFoundException。典型解决方式是:在SpringBoot项目里排除Spark自带的jackson依赖,或者统一指定一个兼容版本。

建议直接用Maven的dependencyManagement把jackson版本锁死,比如统一用2.12.3。另外,Spark作业和SpringBoot应用尽量打成不同的Jar包,用Spark-submit提交作业的是独立的一份,SpringBoot服务里不要打包Spark依赖,这样彻底规避运行时版本冲突。

还有JDK版本。我上面提过,Spark 3.x对Java 17支持还不算完善,Level和反射都用了一些老API。毕设阶段Java 8是最稳的,实在要用Java 11也可以,但别为了新而新。

6.2 数据倾斜与内存溢出,实战中最重要的坑

当行为数据量大了之后,ALS训练阶段非常容易出现OOM或者Executor Lost。这里有个算法原理层面的原因:recommendForAllUsers需要对每个用户做一次“矩阵-向量”计算,如果某个用户的行为数据特别多(比如一个刷单测试账号产生了上万条行为),这个用户所在的分区就比其他分区重太多,形成数据倾斜。

解决数据倾斜我用的三层手段。第一步,在构建评分矩阵时,先做一次行数统计,删除行为数超过阈值(比如500条)的异常用户,这些用户大部分是测试号或爬虫。第二步,调整Spark运行参数,增大单个Task的内存:

spark-submit \ --class com.travel.recommend.RecommendTask \ --master spark://node1:7077 \ --executor-memory 4g \ --driver-memory 2g \ --conf spark.sql.shuffle.partitions=200 \ travel-recommend.jar

第三步,如果自定义算子导致数据倾斜,可以对key增加随机前缀再聚合,做两阶段聚合。但ALS内部的实现是我们改不了的,所以主要靠前两步。

6.3 离线任务调度与重复跑批问题

推荐任务不能手动跑,那太不工程化了。我建议用Quartz或者XXL-Job定时触发。Quartz简单易用,直接集成在SpringBoot里,配置一个Cron表达式:

// 每6小时执行一次推荐任务 @Scheduled(cron = "0 0 */6 * * ?") public void runRecommendTask() { recommendTaskService.execute(); }

但这里有个“重复跑批”的细节:如果上一次任务没跑完,新的定时触发又开始了,两个任务同时读写t_recommend表会造成数据混乱。解决方式是加一个分布式锁,用Redis的setnx做一个任务锁:

boolean locked = redisTemplate.opsForValue().setIfAbsent("lock:recommend", "1", 30, TimeUnit.MINUTES); if (!locked) { log.warn("推荐任务正在执行中,跳过本次触发"); return; }

执行完成后释放锁。这样即使任务本该30分钟完成,机器卡顿导致跑了2小时,锁还不会过期误放。另外,每次跑批前把上一批的t_recommend表数据标记为失效(用批次号逻辑删除),不要直接物理删除,方便回溯对比。

实现时我建议在execute()方法里面还要加幂等控制。任务开始前先查t_recommend_log_task表,看当前批次是否已经跑过,如果跑过则直接跳过。这样即使部署时误触发多次,也不会产生脏数据。

6.4 热点问题:本地开发跑不动Spark怎么办

这是个很现实的问题。很多同学的电脑配置一般,本地跑Spark任务要启动好几个JVM,内存分分钟爆掉。我的建议是:本地开发阶段,把Spark配置成local[*]模式,数据量限制在10万条以内,把逻辑跑通。等需要演示或者做压测时,再部署到云服务器或者实验室的集群上。

如果你没有集群,也可以用单机模式跑Spark,把它当本地计算引擎用。毕设答辩演示本地Demo已经足够,只要你在文档里写清楚“生产环境可扩展为Standalone或YARN集群模式”,理论层面就完整了。

另一个本地开发小技巧是,用H2或者SQLite内存数据库先替代MySQL做测试,避免频繁连接真实数据库占用端口和资源。等逻辑验证通过,再切回MySQL。

6.5 冷启动评估的一个细节

最后再补充一个关于效果评估的坑。冷启动用户通常没有评分历史,所以你在算RMSE时只能评估“有行为的老用户”。如果你想评估整个系统的冷启动质量,一个更好的方式是统计“新用户注册后7天内的平均推荐文章点击率”。刚开始可能只有2%-3%,等你把热门榜、地域推荐、随机探索这些策略都加上之后,能提升到10%以上。这个数字会是你整个毕设最有说服力的一张图表。

结语:做完这个项目的几点真实感悟

我陆陆续续带过不少做推荐系统毕设的同学,这个课题最大的特点就是:门槛不算高,但上限非常高。如果你只是把SpringBoot + Spark搭起来,然后调用一下ALS,那确实是个不错的“实验报告”。但如果你能把整个数据流转链路串起来,从用户行为埋点、离线训练、多路召回、结果缓存、实时扩展,到效果评估,每一层都有清晰的方案和代码,那这个系统就是一个真正完整的工程作品了。

我个人在实际开发中体会最深的,就是“不要把模型当黑盒”。把ALS的参数撑起来、把数据刷起来之后,多看一下每个环节的输出,多统计几版模型的推荐点击率,你会发现推荐结果的差异性比你想象的大。调试这套系统最耗时的地方不是写代码,而是反复调整Rank、RegParam和权重表之间的配合。一旦搭好了这套数据链路,后续加新的推荐策略、新的行为类型,都只是往管道里加新的“料”,整个系统会越来越聪明。

这个方向后续还可以扩展的方向有:接入Flink做实时行为特征计算、用图神经网络做景点关系挖掘、加一个基于大模型的行程规划agent来生成多日游路线。但万变不离其宗,核心的“数据 -> 计算 -> 服务”这条链路,你在做这个毕设的时候已经打通了。这就是这个项目给你留下的最大财富。

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

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

立即咨询