☰
SpringBoot+Spark用户行为分析与推荐系统实战指南
2026/10/5 3:54:52 网站建设 项目流程

先说一下这个选题。SpringBoot 加 Spark 做用户行为分析,这几年在毕设里几乎成了标配组合,但很多同学做到一半就卡住了——要么是环境搭不起来,要么是 Spark 算完的结果不知道怎么优雅地暴露成接口给前端用,要么是埋点数据到手了却不知道清洗成什么样子才算干净。这篇文章把我自己做这类项目时踩过的坑、验证过的方案完整写出来,从架构选型到代码落地,再到排查技巧,希望你看完能少走几趟弯路。

1. 项目拆解与方案选型

1.1 为什么是 SpringBoot + Spark 而不是其他组合

先聊“为什么”。用户行为数据挖掘这个需求本身,拆开看是两层:数据侧要处理海量日志、做清洗、聚合、特征提取、甚至跑算法模型;业务侧要把分析结果变成接口、页面、报表,让管理员能看、能查、能导出。

如果只用单机 Java 技术栈,比如 SpringBoot + MySQL + 定时任务,分析几万条数据没问题,但日志量一旦上来(尤其是埋了点 PV/UV 之后),SQL 聚合慢得让人怀疑人生。如果只用 Spark 而不要 SpringBoot,那你的分析结果只能通过命令行或者临时脚本查看,根本没法做成一个“系统”,答辩的时候演示环节会很尴尬。

所以这个组合的核心逻辑是:Spark 负责“算”,SpringBoot 负责“接”。数据清洗和挖掘全部在 Spark 侧完成后,写入 MySQL、Redis、ES 这类存储,然后 SpringBoot 通过 MyBatis 或 JPA 把它读取出来,以 REST API 的形式提供给 Vue 或 Element UI 做的后台管理界面。整个链路清晰、可拆解也方便讲清楚,非常切合毕业设计的课题要求。

1.2 整体架构与模块边界

我建议把系统拆成下面几个独立模块,各管一摊:

模块核心职责关键技术点
日志采集模块接收前端或客户端上报的行为日志Controller 接收、Kafka 中转(可选)、落盘 HDFS
数据清洗与预处理过滤脏数据、去重、补全字段Spark 算子:filter、dropDuplicates、withColumn
用户画像模块标签化用户特征Spark SQL 聚合、UDF 自定义函数
偏好挖掘模块统计品类/商品/内容偏好groupBy + agg + 窗口函数
推荐算法模块基于协同过滤产出个性化推荐ALS 模型训练、模型持久化和加载预测
结果存储与接口层将结果暴露给前端Redis 缓存、MySQL 持久化、RestTemplate/OpenFeign

模块之间的依赖关系在工程里也要理顺,否则写到最后就是一团浆糊。我的做法是:Spark 作业通过命令行参数驱动,独立成 job 包;SpringBoot 只负责定时触发和结果读取,绝不在业务代码里直接写 Spark 逻辑。这样既方便本地调试 Spark,又不会把 Web 应用的启动时间拖慢——毕竟 SpringBoot 里嵌入 SparkContext 会导致启动变重、显存或内存占用也不好控制。

1.3 毕业设计论文与工程实现的对齐

还有一点值得说:论文怎么写,和工程怎么做,是很多人忽略的“双线对齐”。论文的核心章节一般是:需求分析 → 系统设计 → 数据挖掘模型设计 → 系统实现 → 系统测试。而工程代码的模块划分最好能一一映射到论文章节。比如论文第四章“用户行为分析模型”,代码里就应该能清晰找到UserProfileJob.java、PreferenceAnalysisJob.java、RecommendJob.java这三个类。答辩的时候老师问“你的模型在哪”,你打开 IDE 的项目树直接指给他看,比你翻半天 PDF 有用得多。

2. 环境搭建与 Spark 开发准备

2.1 版本选型——最容易被坑的地方

先强调一个最容易炸的坑:版本不匹配。如果你的 Spark 是 3.x,而 SpringBoot 还是 2.4 这种老版本,问题还不太大;但 Spark 3.3+ 配套的 Scala 2.12/2.13、Hadoop 3.x,以及 JDK 8/11 的兼容关系一定要先确认清楚。我推荐一版我实测过跑通的组合:

组件推荐版本说明
JDK1.8 或 11JDK8 最省心,JDK11 需要额外调整模块权限
SpringBoot2.7.x3.x 也可以,但部分老旧依赖要适配
Spark3.3.x稳定且资料多,对 Scala 2.12 兼容性好
Scala2.12.15和 Spark 配套,不要自己乱升级
Hadoop3.3.4仅本地模式可省略 HDFS 配置
MySQL8.x注意驱动用com.mysql.cj.jdbc.Driver
Maven3.8+构建管理

版本选型为什么重要?举个例子,很多人下载了 Spark 3.4,然后发现它默认用 Scala 2.13,你本地 Maven 里又引入 2.12 的包,运行时报java.lang.NoSuchMethodError,找错能找两天。另外,MySQL 驱动如果还用旧的com.mysql.jdbc.Driver,在 SpringBoot 2.7 下会直接启动失败,报No suitable driver found,也容易被误导成数据库连接问题。

2.2 Maven 工程如何组织 Spark 依赖

Spark 依赖的坐标和普通 JAR 不一样,它的 artifactId 里有 Scala 版本号。我习惯这样配置:

<properties> <spark.version>3.3.2</spark.version> <scala.version>2.12</scala.version> </properties> <dependencies> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_${scala.version}</artifactId> <version>${spark.version}</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_${scala.version}</artifactId> <version>${spark.version}</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-mllib_${scala.version}</artifactId> <version>${spark.version}</version> </dependency> </dependencies>

这里有一个容易被忽略的细节:spark-mllib里的 ALS 推荐算法在spark-sql之外还要单独引入。如果你只做普通聚合统计,不写协同过滤,spark-core加上spark-sql就够了;但涉及模型训练,一定要带上 MLLib。依赖传递的问题也很多——举例说,Spark 的日志框架和 SpringBoot 的 Logback 会冲突,运行时经典表现是SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder"。我的做法是排除掉 Spark 自带的slf4j-log4j12,统一用 Logback。

2.3 本地模式配置与集群模式切换

开发阶段强烈建议用本地模式,也就是spark.master=local[*]。你可以在 SpringBoot 的application.yml里配置一个 profile 来切换:

spark: app-name: user-behavior-analysis master: local[*] executor-memory: 2g driver-memory: 1g

然后写一个SparkConfig配置类:

@Configuration public class SparkConfig { @Value("${spark.master}") private String master; @Value("${spark.app-name}") private String appName; @Bean public SparkSession sparkSession() { return SparkSession.builder() .appName(appName) .master(master) .config("spark.sql.shuffle.partitions", "4") .config("spark.sql.adaptive.enabled", "true") .getOrCreate(); } }

开发机上用local[*],它的意思是“使用本机所有核心并行跑”。这个模式的好处是不用启动 HDFS、不用管 YARN,直接就能跑通逻辑。到了演示或者生产环境,你只需要把master改成yarn,配合spark-submit提交即可,代码不用大改。不过必须提醒你:SpringBoot 启动时如果直接创建 SparkSession,每次启动都会默认打印大量 Info 日志,感官上会误以为应用卡住了,建议把日志级别调成 WARN:

sparkSession.sparkContext().setLogLevel("WARN");

3. 用户行为数据模型设计与预处理

3.1 行为日志埋点字段设计

要分析用户行为,首先得定义“行为”有哪些。我常用的行为类型包括:

  • 浏览商品 / 访问页面:view
  • 搜索关键词:search
  • 收藏商品:favorite
  • 加入购物车:cart
  • 提交订单:order
  • 支付成功:pay

每条行为日志建议包含这些字段:

字段名示例值说明
user_id100234用户唯一标识,匿名用户置为 -1
session_id8fdc43a3一次会话 ID
item_id / product_idP1024被操作对象 ID
behaviorview行为类型枚举
category_idC22可选,用于品类偏好分析
from_channelapp / pc / wechat来源渠道
timestamp1712483200000毫秒时间戳
stay_time45停留秒数,浏览页才有

这里要特别提醒:埋点越规范,后面分析越省力。如果behavior字段习惯用中文(浏览、点击),Spark 处理时要用 UDF 做映射,徒增工作量。更好的做法是埋点时就写枚举值字母,前端或者小程序端展示层再映射成中文,两边都干净。

3.2 SpringBoot 如何接收高并发日志上报

通常的做法是提供一个/api/log/collect接口,前端用navigator.sendBeacon或fetch在用户产生行为时上报一条 JSON。后端接口本身很简单,就是接收数据、异步转发到消息队列或直接落库。但有一个容易忽略的问题:如果每个行为都直接写数据库,流量稍微一大数据库就会被写穿。我的方案有两个可选:

  1. 先写入 Redis 的 List 结构,批量化落库。
  2. 使用 Kafka 作为缓冲层,Spark 用 Streaming 模式消费。

如果你毕设不想引入 Kafka,可以用方案 1,代码大致:

@RestController @RequestMapping("/api/log") public class LogCollectController { @Autowired private StringRedisTemplate redisTemplate; @PostMapping("/collect") public Result collect(@RequestBody UserBehaviorLog log) { // 同步校验日志是否合法 if (log.getUserId() == null || log.getBehavior() == null) { return Result.error("参数缺失"); } // 异步写入 Redis 列表,后面由定时任务批量刷入 redisTemplate.opsForList().leftPush( "behavior:logs", JSON.toJSONString(log) ); return Result.success(); } }

然后再写一个定时任务或者直接写一段 Spark 批量任务读取 Redis 里的数据。更进一步,可以直接用 Spark 对接 Redis 数据源,但那样需要额外引入 Redis 连接工具的转换逻辑。为了控制复杂度,我更建议定时把 Redis 里攒的日志写回 HDFS 或 MySQL 临时表,再用 Spark 分析。

3.3 数据清洗与 ETL 的 Spark 实现

日志数据本身是很脏的。清洗逻辑我总结了四个常规步骤:

  • 字段缺失处理:user_id为空的丢弃,item_id为空的如果是浏览记录则可以保留但置为unknown。
  • 时间字段标准化:统一转成 yyyy-MM-dd HH:mm:ss 格式,方便后续按小时/天聚合。
  • 去重:同一用户在同一秒内对同一个商品的行为,可以认为是重复点击,用dropDuplicates处理。
  • 行为合法性校验:行为类型不在枚举范围里的直接过滤。

一个典型的 Spark 清洗代码如下:

import org.apache.spark.sql.{SparkSession, DataFrame} import org.apache.spark.sql.functions._ object LogCleanJob { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("LogCleanJob") .master("local[*]") .getOrCreate() import spark.implicits._ val rawDF: DataFrame = spark.read.json("/data/behavior_logs") val validBehaviors = Seq("view", "search", "favorite", "cart", "order", "pay") val cleanedDF = rawDF .filter($"user_id".isNotNull && $"behavior".isNotNull) .filter($"behavior".isin(validBehaviors: _*)) .withColumn("timestamp", to_timestamp($"timestamp" / 1000)) .dropDuplicates("user_id", "item_id", "behavior", "timestamp") .withColumn("date", date_format($"timestamp", "yyyy-MM-dd")) cleanedDF.write.mode("overwrite").parquet("/data/behavior_cleaned") spark.stop() } }

这段代码要解释一个关键点:去重字段组合。只对“用户+商品+行为+时间”四元组去重,才能避免把同一用户真实浏览不同商品的情况误删。很多新手用dropDuplicates("user_id")去重,那会直接把一个用户的所有行为删得只剩一条,后面的分析基本全废。还有,JSON 读取时时间戳常常是 Long 类型毫秒值,Spark 的to_timestamp函数接收的是秒,所以要先除以 1000,否则结果是 1970 年。

我这里用了 Scala 写 Spark 作业,因为 Spark 的 Scala API 在写复杂处理链时确实更简洁。如果你只熟悉 Java,用 Spark 的 Java API 也可以写,但代码要长三分之一以上,而且很多函数式写法在 Java 里会显得繁琐。

3.4 清洗后的质量评估

清洗完成不要急着往下走,先肉眼看一下统计结果。我一般会输出一个质量报告:总日志量、有效日志量、去重后日志量、各行为类型占比、活跃用户数。这些数字既能验证清洗逻辑是否正确,同时也是论文里“系统测试”章节的素材。说个真实经历:我第一版清洗代码把时间戳按秒处理,结果日期全变成了 1970 年的某几天,如果不是输出日活统计时发现“每天百万用户”这种诡异数据,根本察觉不到是时间戳单位的问题。

4. 用户画像与偏好分析实现

4.1 基于 Spark SQL 的聚合统计

用户画像是用户行为分析的“第一张牌”。简单来说,就是要回答:这个用户是谁?他喜欢什么?他活跃在什么时间段?我通常先用 Spark SQL 做常规聚合:

SELECT user_id, COUNT(CASE WHEN behavior = 'view' THEN 1 END) AS view_cnt, COUNT(CASE WHEN behavior = 'cart' THEN 1 END) AS cart_cnt, COUNT(CASE WHEN behavior = 'order' THEN 1 END) AS order_cnt, COUNT(DISTINCT category_id) AS category_cnt, COUNT(DISTINCT item_id) AS item_cnt, MAX(CASE WHEN behavior = 'order' THEN category_id END) AS last_order_category FROM behavior_cleaned GROUP BY user_id

这段 SQL 在 Spark 里可以直接通过spark.sql()执行。它产出的结果表基本就是用户特征表的雏形。在此基础上,我们可以根据order_cnt给用户打上“高消费”、“普通”、“潜水”的标签;根据时段偏好把用户分为“夜猫子型”、“上班族型”等。打标签的规则在论文里可以写得有逻辑:阈值分箱(如订单数 > 5 为高活跃)比拍脑袋定性要好讲得多。

4.2 自定义 UDF 构建用户标签体系

有了基础聚合数据,接下来就是“标签化”。Spark 的自定义函数(UDF)在这里派上用场。比如我要根据活跃度打分:

val activeLevel = udf((view: Int, cart: Int, order: Int) => { val score = view * 1 + cart * 2 + order * 5 if (score >= 100) "核心用户" else if (score >= 30) "活跃用户" else if (score >= 5) "普通用户" else "沉默用户" }) val profileDF = aggDF.withColumn("active_level", activeLevel($"view_cnt", $"cart_cnt", $"order_cnt"))

这类规则型的标签好处是可解释性很强,答辩时老师问“为什么给他打这个标签”,你可以直接说“因为最近 30 天下单 6 次、浏览 120 次,加权得分超过阈值,规则明确”。如果用到机器学习聚类打分,反而很难几句话解释清楚,对毕设来说不划算。当然,如果你的课题名称里有“智能”两个字,想加一点机器学习的料,聚类也不是不行,但规则标签要保留,作为对照。

4.3 偏好分析:从行为到品类偏好度

品类偏好分析是另一个常见模块。核心思路是:用户对某个品类的偏好度 = 用户在该品类的有效行为数 / 用户全品类有效行为数。更严谨一点可以引入时间衰减因子——越近的行为权重越高。用 Spark 实现时间衰减可以用 UDF:weight = 1 / (days_ago + 1)。逻辑如下:

import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val prefDF = cleanedDF .groupBy("user_id", "category_id") .agg( sum( when(col("behavior").isin("cart", "order"), 3) .when(col("behavior") === "view", 1) .otherwise(2) ).as("raw_score") ) val userTotal = prefDF .groupBy("user_id") .agg(sum("raw_score").as("total_score")) val ratioDF = prefDF.join(userTotal, "user_id") .withColumn("pref_ratio", col("raw_score") / col("total_score"))

这里的行为赋分逻辑也是可以写进论文作为“偏好度模型”的。顺便说一句,窗口函数在 Spark 里做排行非常好用,比如每个用户浏览最多的 Top3 品类:

val topCategory = prefDF.withColumn( "rank", row_number().over(Window.partitionBy("user_id").orderBy(col("raw_score").desc)) ).filter(col("rank") <= 3)

窗口函数是 Spark 分析中必考的知识点,能在代码里体现这个概念,答辩被问到的概率极高,所以不要回避它。

4.4 结果存储结构设计

分析完的数据要有地方放。我的建议是建这三张表:

  • user_profile:用户基础画像标签表,按天更新覆盖。
  • user_category_pref:用户-品类偏好分数表。
  • item_view_rank:商品浏览热度榜,供前端展示热门商品。

存储格式上,大结果集写成 Parquet 存 HDFS,但要给 SpringBoot 接口用,还得同步一份到 MySQL;小结果集直接写 Redis 可以显著降低接口延迟。实践中我一般是 Spark 算完之后,通过 JDBC 把结果写 MySQL。Spark 写 MySQL 的姿势要稍微注意下,直接df.write.jdbc()在大数据量下会很慢,而且容易把连接池打满。更稳的方式是先把结果 collect 成少量分区,再 foreachPartition 分批插入:

prefDF.write .mode("overwrite") .option("truncate", "true") .jdbc("jdbc:mysql://localhost:3306/user_analysis?useSSL=false&serverTimezone=Asia/Shanghai", "user_category_pref", connectionProperties)

如果你是在本地模式测试,数据量不大,用jdbc直写也没问题。但如果数据到了百万级以上,强烈建议用foreachPartition加批量 insert,否则会有连接超时报错。

5. 基于协同过滤的推荐模块

5.1 ALS 模型的原理与在 Spark 中的落地

用户行为数据挖掘的进阶内容是推荐系统。这儿我用的是 Spark MLLib 里的 ALS(交替最小二乘法)。它的核心思想很简单:把“用户对商品的偏好矩阵”分解成两个低维矩阵的乘积——一个代表用户特征,一个代表商品特征,然后通过内积预测用户对未购买商品的评分。

代码示例:

import org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.functions._ val ratings = cleanedDF .filter(col("behavior").isin("view", "cart", "order", "pay")) .groupBy("user_id", "item_id") .agg( sum( when(col("behavior") === "pay", 5) .when(col("behavior") === "order", 4) .when(col("behavior") === "cart", 3) .otherwise(1) ).as("rating") ) val als = new ALS() .setMaxIter(10) .setRegParam(0.1) .setUserCol("user_id") .setItemCol("item_id") .setRatingCol("rating") .setColdStartStrategy("drop") val model = als.fit(ratings) model.write.save("/models/als_model")

有几个点必须注意:

  • 评分映射逻辑:把pay映射为 5、order映射为 4、cart映射为 3、view映射为 1,这是一种启发式做法。你可以根据自己的需求调整权重。
  • ColdStartStrategy 一定要设成 drop:否则遇到新用户或新商品,ALS 预测时会产生 NaN 评分,接口返回一堆 null,前端展示直接崩。
  • 模型保存后,SpringBoot 读模型做预测是另一套逻辑。最简单的做法是:Spark 离线算完,直接把推荐结果写到 Redis 或 MySQL,SpringBoot 只做读操作。这样就不用在 Web 进程里维护 SparkSession。

5.2 推荐结果的召回与排序

协同过滤算完后,每个用户会产生一个包含若干商品 ID 和评分的列表。不能直接把 Top10 给前端,还需要做一层过滤和排序:

  • 过滤用户已购买过的商品。
  • 过滤下架商品(可以 join 商品表做过滤)。
  • 与热门榜做一定比例的融合:比如 60% 个性化结果 + 40% 热门结果,避免冷启动用户看到空列表。

融合的代码逻辑在 SpringBoot 里做即可:

public List<RecommendItem> getRecommendItems(Long userId, int size) { // 1. 从 Redis 获取协同过滤结果 List<RecommendItem> personal = recommendMapper.selectByUserId(userId); // 2. 过滤已购买 List<Long> boughtIds = orderMapper.selectItemIdsByUser(userId); personal.removeIf(item -> boughtIds.contains(item.getItemId())); // 3. 若不足则用热门补足 if (personal.size() < size) { List<RecommendItem> hotItems = itemMapper.selectHotItems(size - personal.size()); personal.addAll(hotItems); } return personal.subList(0, Math.min(size, personal.size())); }

这段代码在答辩时可以清晰地讲出“推荐结果不是模型出来就完事了,还得考虑业务规则”。这是加分项。很多同学模型跑出来,接口一返回一堆空推荐就不知所措,其实关键就是没有做冷启动兜底。

5.3 离线计算与定时调度

整个分析流程的调度也值得说说。我用的是 SpringBoot 自带的@Scheduled注解,每天晚上凌晨跑一次全量分析任务:

@Component public class AnalysisScheduler { @Scheduled(cron = "0 0 2 * * ?") public void runDailyAnalysis() { // 调用 Spark 提交逻辑,或用 ProcessBuilder 调用 spark-submit } }

一种更干净的实现是,在工程里把 Spark 任务打成可执行 JAR,然后用ProcessBuilder去触发spark-submit。SpringBoot 只负责记录任务状态、开始时间和结束时间。这样 Web 应用本身不会因为 Spark Job 异常崩溃。还有一点,Spark 任务跑完后要主动调用spark.stop()释放资源,不然第二天定时任务再启动时会报端口冲突或内存不足。

6. 常见问题与排查技巧实录

6.1 Spark 作业频繁 OOM

这是本地模式最常见的头号 bug。OOM 的原因,一半是数据倾斜,一半是分区数太少。排查步骤是:看 Spark UI 里的 Stage 明细,如果某个 Task 处理的数据量是其他 Task 的好几倍,基本就是数据倾斜了。解决方案:

  • 提高分区数:repartition(col("user_id"), 200)或设置spark.sql.shuffle.partitions=200。
  • 对 key 加盐,做两阶段聚合。
  • 给 driver 和 executor 加大内存:本地测试可以设driver-memory=4g。

新手容易忽略的一点是:本地跑 Spark 时,SpringBoot 本身也占内存,如果你 IDE 里的 VM 选项只给了 512m,Spark 一启动就 GG。建议把 IDE 的运行内存调到 2g 以上。

6.2 SpringBoot 启动被 Spark 拖慢

问题描述:SpringBoot 启动要 2 分钟,日志刷屏。原因通常是 SparkSession 被@Bean初始化并且打印了大量 INFO。解法:

  • 设置setLogLevel("WARN")
  • 把 SparkSession 定义为懒加载,只有真正调用分析接口时才创建。
  • 或者在单独的 Maven Profile 里把 Spark 隔离,本地开发 Web 时不加载 Spark 相关依赖。

如果用了懒加载,要注意一点:第一次请求分析接口时,创建 SparkSession 需要约 10 秒左右,前端要设置超时时间,并且最好有 loading 状态,否则用户会以为接口挂了。

6.3 MySQL 写入中文乱码和数据截断

出现这个问题的原因一般是连接串没带字符集参数。修复很简单:JDBC 连接串里加上useUnicode=true&characterEncoding=utf8,同时确保 MySQL 表本身是 utf8mb4 字符集。如果你在 Spark 里用write.jdbc往 MySQL 写数据,建议在连接属性里也显式指定:

user=root password=123456 useUnicode=true characterEncoding=utf8

还有一个非常隐蔽的坑:MySQL 的tinyint在某些驱动版本下会被自动映射成Boolean,导致你往 Java 实体里塞值时出现ClassCastException。解决方法是调整 JDBC URL 里的tinyInt1isBit=false,或者在实体里用Integer接收然后做转换。

6.4 推荐接口返回空列表

这个问题出现时不要慌。依次排查:Redis 里有没有对应 key?模型训练集是否覆盖了这个用户?冷启动策略是否生效?如果用户是个新注册的号,没有历史行为,ALS 根本不可能预测出结果。你的兜底方案就要把热门榜顶上来。这也是为什么我在 5.2 里强调“协同过滤 + 热门内容融合”,不是理论空谈,而是真的会遇到。

6.5 数据时间显示 1970-01-01

这个坑前文已经提及,但特别值得单独列为一条。所有时间戳处理前,先确认它是秒还是毫秒。Spark 的from_unixtime默认处理秒,Java 的System.currentTimeMillis()返回毫秒。调试办法很简单:随便拿一个时间戳去在线转换工具验证下,看是不是对的,不要凭感觉。

7. 从毕设项目到真实工程的进阶思考

做到这里,一个完整的 SpringBoot + Spark 用户行为分析系统已经能跑了。但如果你想让这个项目在答辩中更出彩,还有几个可以“免费加分”的点:

第一,引入 Redis 做实时热榜。前端展示的“热门商品 Top10”尽量不要每次都查 MySQL,而是让 Spark 算完后写入 Redis 的 ZSet,接口直接ZREVRANGE取前 N 名,响应时间能压到几十毫秒。这个设计在答辩的时候说“用 Redis 做缓存加速读热点数据”,很加分。

第二,增加可视化分析页。很多人只在后端做了接口,前端却只是表格展示。如果加上 ECharts 的折线图、饼图、热力图,展示用户活跃时段分布、行为类型占比和品类偏好分布,整个系统立刻显得完整很多。

第三,设计一个数据质量监控面板。统计每日日志量、清洗率、异常率,出现异常时告警。这个模块代码量不大,但“数据质量”四个字在导师和评委眼里非常专业。

第四,考虑把数据源从 MySQL 换到更贴近生产环境的方案。本地日志文件虽然简单,但如果在项目描述里写“对接 Kafka 实时消息流”,会更有“大数据实时处理”的味道。SpringBoot 里用 Kafka 客户端作为生产者,Spark 用readStream消费,虽然开发量上了一个台阶,但这也是一个完整的实时数仓雏形。

我个人从这类项目里最大的体会是:技术栈永远是工具,真正拉开差距的是数据建模的思路和排查问题的能力。很多人卡在“Spark 跑通了但不知道下一步干什么”,本质是没把业务问题拆解成可以用数据回答的问题。你问“用户喜欢什么”,对应的就是“品类偏好分数”;你问“给这个用户推荐什么”,对应的就是“ALS 评分 + 热度兜底”。把这句话想透,你的系统设计就不会偏。

最后再分享一个小技巧:写完 Spark 作业后,一定把执行计划里 Shuffle 的数据量看一眼。Shuffle 数据越大,作业越慢。很多所谓“性能优化”问题,其实只是多一个filter、少一个join的事。你在自己机器上把作业从 2 分钟压到 30 秒,这个成就感,比单纯把功能跑通要实在得多。

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

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

立即咨询