1. 项目概述与整体架构设计
1.1 这个项目到底解决什么问题
做这个课题之前,我先去市场上逛了一圈。宠物用品的价格在不同平台确实能差出不少,同款猫粮这边标价260,那边做活动能低到190,用户想买又怕买贵,只能一家一家去翻。但真正推动我做这个项目的,是另一个更实际的问题:当商品数据规模达到几十万条甚至上百万条时,靠人力一条条去比价、去筛选、去人工推荐,已经完全不可能了。传统的关系型数据库单表查个价格索引撑死几十毫秒,一旦要实时比对全网商品,再叠加每个用户的行为偏好计算,很快就卡死。
所以我这个项目的定位很清晰:做一个基于大数据技术栈的宠物商品信息比价及推荐系统。核心能力拆开就两块,第一块是基于多数据来源的宠物商品采集、清洗、价格对比,第二块是根据用户浏览、收藏、购买行为做个性化推荐。两块能力都跑在Hadoop和Spark这套大数据底座上,业务接口和后台管理用SpringBoot来做,最后配一块可视化大屏把商品价格分布、推荐效果、销量趋势这些关键指标实时展示出来。
系统适合谁参考?一个是做毕业设计或者课程设计的学生,另一个是刚入门大数据、想看看Hadoop+Spark+SpringBoot怎么整合的真实项目的开发者。这套系统里面涉及的组件很多,从采集端到计算端到展示端都有,能帮你把学校学的零散知识点串成一条线。
1.2 为什么选这四件套
技术选型这块我纠结过一段时间。当时摆在面前的有两条路:一条是老老实实用MySQL+Redis+SpringBoot,开发效率高、工期短,但本质还是一个传统CRUD项目,撑不起“大数据”三个字;另一条就是用分布式存储和分布式计算,技术含量上去了,复杂度也跟着上去了。
最后我选择的是Hadoop作为底层分布式存储,Spark作为计算引擎,SpringBoot作为业务接口层,Vue+ECharts作为可视化展示层。这套组合的好处在于:
Hadoop的HDFS在处理海量商品快照数据时有天然优势。宠物商品的比价系统需要周期性抓取各平台的价格快照,一份快照可能上千万条记录,全部扔进MySQL会导致单表爆炸,HDFS这边用分区目录按天存,配合Hive元数据管理,查询和归档都很顺手。
Spark替代MapReduce做离线计算是必然选择。同样做用户行为数据的聚合分析,MapReduce每个Job跑一个多小时是常态,Spark基于内存的DAG计算能把时间压缩到一个零头。推荐引擎里面最关键的相似度矩阵计算,用Spark的MLlib库跑,效率和代码量都比手写MapReduce友好太多。
SpringBoot负责把后台能力和用户侧衔接起来。前端收到用户的行为数据,打到SpringBoot的Controller,Controller写入Kafka主题或者调用Spark接口触发计算;当Spark算完推荐结果,落回MySQL或者Redis,前端再从SpringBoot读出来展示。
可视化大屏是对外展示和答辩的“门面”。这个系统如果只有几个API接口,那看起来就是一个纯后端项目;但加上大屏,就是一个闭环的完整应用,而且能把Spark算出来的数据价值直观呈现出来。
从整体架构上看:
| 层级 | 技术栈 | 职责 |
|---|---|---|
| 数据接入层 | 爬虫脚本 + Kafka | 采集各平台宠物商品价格、用户行为数据 |
| 存储层 | HDFS + Hive + MySQL + Redis | 海量快照存储、元数据管理、业务数据落地 |
| 计算层 | Spark核心 + Spark SQL + MLlib | 价格聚合、相似度计算、推荐列表生成 |
| 业务层 | SpringBoot | 用户管理、商品检索、比价接口、推荐接口 |
| 展示层 | Vue + ECharts | 商品详情、价格曲线、可视化大屏 |
这个分层结构让每个组件都聚焦在自己的职责上,出了问题也能快速定位。
1.3 完整的数据链路串讲
让我把一条用户请求的完整链路走一遍,这样理解起来更清晰。用户在Web端打开商品页,点击“查看历史价格曲线”,前端请求打到SpringBoot的价格查询接口,这个接口首先查Redis缓存,缓存没有再去查MySQL的最近快照表,如果还要更长的历史区间,就通过HiveSQL去查HDFS上的月度分区数据。
用户浏览了几件商品后,在页面停顿了一会儿,前端把这次浏览行为异步上报到Kafka的Topic,Spark Streaming任务每隔几十秒从Kafka拉一次行为数据,更新用户画像,并触发增量式的相似商品推荐计算。计算结果写回Redis的推荐桶,用户下次刷新页面时直接从SpringBoot读到推荐列表,整个过程对用户来说几乎是无感知的。
再说比价流程。采集程序定时把三家主流电商平台的同款宠物商品链接、价格、优惠信息拉下来,落到HDFS的同一份商品维度表里,Spark SQL按商品ID做分组聚合,计算最低价、最高价、均价,最后把聚合结果写回MySQL,供前端直接展示。这套流程听起来不复杂,但真正做的时候,数据清洗、去重、单位归一化这些坑我一个一个踩了个遍,后面单独开一节讲。
2. Hadoop与Spark环境搭建实战记录
2.1 Hadoop伪分布式安装与配置
环境这块我是在Ubuntu 20.04上装的,用的Hadoop 3.3.4版本。之前很多教程还在用2.x,我建议直接用3.x,一方面HDFS支持了纠删码,另一方面和Spark 3.x的兼容性更好。
安装流程说穿了就四步:配JDK、配SSH免密、改配置文件、启动验证。JDK我装的JDK 8,因为Spark 3.2以下的版本对JDK 11的支持还不完整,与其后面调各种兼容问题,不如老老实实用8。
SSH免密这个步骤容易栽跟头,很多人第一次做集群根目录都搞不对。执行ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa生成密钥后,要把公钥追加到authorized_keys里:
cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys注意authorized_keys的权限一定要是600,~/.ssh目录权限必须是700,权限不对直接不生效,而且报错信息还不明显,只提示Permission denied。
核心配置文件就三个:core-site.xml、hdfs-site.xml、yarn-site.xml。我当时的配置可以参考:
<!-- core-site.xml --> <configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/home/bigdata/hadoop_tmp</value> </property> </configuration><!-- hdfs-site.xml --> <configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/home/bigdata/hadoop_tmp/namenode</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/home/bigdata/hadoop_tmp/datanode</value> </property> </configuration>伪分布式模式下dfs.replication必须设为1,否则副本数超过节点数,DataNode会一直报副本不足的告警。这个问题当年让我卡了整整一个下午,看到的全是Replica placement policy异常。
格式化NameNode之前要确认hadoop_tmp目录是空的,格式化命令是hdfs namenode -format,看到“SHUTDOWN_MSG”出现,说明格式化成功,但它只是初始化了元数据存储目录,并不代表集群已经跑起来。启动的时候用sbin/start-dfs.sh,然后jps查看进程,能看到NameNode、DataNode、SecondaryNameNode三个Java进程,才算成功。
2.2 Spark集群部署与任务提交要点
Spark我装的是3.2.1版本,与Hadoop 3.3.4搭配没遇到过兼容性问题。部署模式用的是Standalone,配置上比YARN模式简单,调试也更直观。因为机器资源有限,我开了两台虚拟机搭了一个最小集群,一台Master一台Worker。
Spark安装的关键在于环境变量的配置。在~/.bashrc里面加入:
export SPARK_HOME=/home/bigdata/spark-3.2.1 export PATH=$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin然后修改spark-env.sh指定Java环境和Master地址:
export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64 export SPARK_MASTER_HOST=bigdata-master export SPARK_WORKER_CORES=2 export SPARK_WORKER_MEMORY=4g这里SPARK_WORKER_MEMORY决定了每个Worker节点的可用内存,要根据你自己的物理内存来设,别上来就填8g,启动了半天Worker直接OOM退出,到时候排查起来又得绕一圈。
提交推荐任务的时候,我习惯先把历史行为数据从HDFS加载出来做离线推荐计算,命令大致长这样:
spark-submit \ --master spark://bigdata-master:7077 \ --class com.petmall.recommend.OfflineRecommender \ --executor-memory 2g \ --num-executors 2 \ pet-recommend-1.0.jar--executor-memory和--num-executors这两个参数是调优重点。我一开始给了每个executor 4g内存,结果集群两个Worker一共才8g,任务提交后频繁GC,日志里全是WARN MemoryStore: Not enough space to cache,调成2g反而稳了。Spark的资源估算没那么想当然,Executor内存不是越多越好,要看你整个集群的总资源和并发任务数来平衡。
2.3 数据采集与数仓分层设计
数据是整个系统的血液,采集层我用的Python写爬虫,抓取三家电商平台的宠物猫粮、狗粮、猫砂、宠物玩具四个类目的商品信息。字段包括商品ID、商品标题、分类、品牌、规格、价格、原价、销量、店铺名、抓取时间。
存储这块我借鉴了数据仓库的分层思路,虽然规模比不上企业级数仓,但分层的思想是通用的。ODS层放原始抓取数据,落地成Parquet格式存到HDFS,按日期分区;DWD层做清洗和去重,统一商品名、处理缺失价格、把不同平台的单位统一成克和毫升;ADS层做最后的聚合,算出每个商品在各平台的价格对比、价格波动、推荐指数。
Parquet相比普通文本格式的优势在Spark SQL里体现得很明显。同样是200万条商品快照数据,CSV格式跑一个GROUP BY聚合要40秒,换成Parquet只要12秒左右,而且HDFS上的空间占用少了差不多一半。强烈建议写Spark作业时优先考虑列式存储格式。
3. Spark为核心的推荐模块设计
3.1 推荐算法选型:为什么用协同过滤
宠物商品推荐主要面临一个场景:用户访问商品页,系统要给出一组“猜你喜欢”的列表。我当时在基于内容的推荐和协同过滤之间选了半天。基于内容适合新商品冷启动,但宠物用品之间不像服装那样有明确的风格标签,你很难用文本特征把“鸡肉味猫粮”和“鱼肉味猫粮”的语义距离算准。
协同过滤的优势在于它只依赖用户行为矩阵,不需要对商品内容做特征工程。用户有浏览、收藏、加购、购买四种行为,我通过加权把行为转成评分。权重设置我用了经验值:浏览1分、收藏3分、加购4分、购买5分,最后再做一个归一化处理,得到用户-商品评分矩阵。
算法层面我选了基于物品的协同过滤,也就是Item-CF。选它的原因很实际:宠物商品的数量级在几十万,小于活跃用户的量级,计算物品相似度矩阵比计算用户相似度矩阵更划算,而且Item-CF的解释性更好——“因为你看了XX猫粮,所以推荐你相似的YY猫粮”这个逻辑用户一眼就能理解。
3.2 Spark MLlib实现相似度计算
Spark实现Item-CF的核心步骤是两段:先构建用户对物品的评分矩阵,再计算物品之间的余弦相似度。
import org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.functions._ // 读取去重后的行为数据 val behaviorDF = spark.sql("SELECT user_id, product_id, score FROM dwd_user_product_score") // ALS 隐式反馈训练 val als = new ALS() .setRank(10) .setMaxIter(10) .setRegParam(0.01) .setUserCol("user_id") .setItemCol("product_id") .setRatingCol("score") .setColdStartStrategy("drop") val model = als.fit(behaviorDF)这里需要注意setColdStartStrategy("drop")这行。如果不设置,模型在预测新用户或新商品时会产生NaN评分,导致后续的推荐列表全是空值。这个坑在评测的时候特别容易踩,因为测试集里的用户往往不在训练集里。
ALS训练完,Spark会输出每个用户的商品因子矩阵和每个商品的特征向量。再通过矩阵乘积算出用户对所有商品的预测评分,取Top N写入Redis:
val userRecs = model.recommendForAllUsers(20) userRecs.write.mode("overwrite").parquet("/recommend/result/user_recs")离线推荐的结果我每天晚上定时跑一次,生成全量用户的Top20推荐列表。白天用户在线的实时推荐,则基于最近一小时的行为增量触发,用Spark Streaming消费Kafka里的事件流,更新特定用户的推荐队列。
3.3 比价计算的实现细节
比价模块本质是一个多维度聚合计算。从不同平台抓回来同款商品后,首先要解决“同款”的识别问题。各平台的标题表述差异极大,比如“金素丽高宠物猫粮 鸡肉味 4磅”和“金装素丽高猫粮 鸡肉 1.81kg”,其实是同一个商品。
我用的是“品牌+系列+规格”三级匹配方案。先通过品牌词表映射出品牌,再提取系列关键词,最后把规格统一换算成克或毫升。匹配逻辑写在Spark的UDF里,计算时先按品牌分组,再在组内用编辑距离和关键词共现来判断相似度,避免全量两两比对带来的笛卡尔积爆炸。
注意,这里有一个很重要的经验:比价过程不要用MySQL的JOIN做全表匹配,我第一版这么写过,200万单表数据跑一个LIKE匹配直接让数据库CPU飙到100%。把数据拉到Spark里做分布式计算,同样的逻辑在集群上3分钟跑完,而且完全不阻塞线上业务。
价格聚合结果写到MySQL表commodity_price_aggregation,字段包括商品ID、平台、最低价、最高价、均价、更新时间。前端展示价格曲线时,用Spark SQL从HDFS上的Parquet快照表按日期查询历史价格,返回给ECharts画折线图。
4. SpringBoot业务层与接口开发
4.1 工程结构与模块划分
SpringBoot这块我用的是2.7.x版本,和JDK8搭配相当稳定。项目采用Maven多模块结构:
pet-system/ ├── pet-common # 公共模块:工具类、统一返回体 ├── pet-gateway # 接口模块:登录、用户管理 ├── pet-product # 商品模块:商品搜索、详情、比价 ├── pet-recommend # 推荐模块:推荐列表、反馈收集 ├── pet-admin # 后台管理:商品管理、价格规则配置 └── pet-dashboard # 大屏数据接口模块拆分的原则是各模块之间尽量解耦,比如pet-recommend只依赖pet-common和数据库,不直接依赖pet-product的表,需要商品信息时通过Feign接口调用,虽然费一点事,但后续扩展时不容易改一处崩一片。
application.yml里的关键配置:
spring: datasource: url: jdbc:mysql://localhost:3306/pet_mall?useUnicode=true&characterEncoding=utf8 username: root password: 123456 redis: host: localhost port: 6379 database: 0 mybatis-plus: mapper-locations: classpath:mapper/*.xml configuration: log-impl: org.apache.ibatis.logging.stdout.StdOutImplRedis在这里承担了两个职责:一个是缓存商品详情和比价结果,减少重复查询数据库的压力;另一个是保存推荐列表。推荐列表的Key我按用户维度设计成rec:user:{userId},Value用JSON数组,过期时间24小时。这样设计的好处是离线推荐任务写Redis,在线接口读Redis,两边不用直接通信。
4.2 比价接口与推荐接口的实现
比价接口的URL设计为POST /api/product/price/compare,参数是商品ID列表。接口逻辑分三步:先查Redis缓存,如果商品ID在缓存中有完整的价格对比数据,直接返回;缓存没有的,去MySQL查聚合表;MySQL也没有历史数据的,再触发一次Spark批处理补充计算。
推荐接口是GET /api/recommend/products?userId=123&page=1&size=10。逻辑很简单:从Redis取推荐列表,如果用户是新用户没有行为数据,就回退到全局热门商品列表。这里要设计好回退机制,否则新用户打开页面永远是一片空白,体验太差。
@Service public class RecommendServiceImpl implements RecommendService { @Autowired private StringRedisTemplate redisTemplate; @Override public List<ProductVO> getRecommendList(Long userId, int page, int size) { // 优先读Redis中的个性推荐 String cacheKey = "rec:user:" + userId; String cached = redisTemplate.opsForValue().get(cacheKey); if (StringUtils.hasText(cached)) { return JSON.parseArray(cached, ProductVO.class) .stream().skip((page - 1) * size).limit(size).collect(Collectors.toList()); } // 缓存没有则走热门兜底 return getHotProducts(page, size); } }这个接口看起来简单,但有个细节容易被忽略:Redis里存的是整个推荐列表,分页要在内存里做,而不是去数据库分页。因为Redis里的数据已经是最新算好的推荐结果,再回数据库查一次反而多此一举,而且数据库里的推荐字段更新不及时。
4.3 前端行为埋点上报
推荐系统要产生效果,必须有用户行为数据喂进来。我在前端页面做了埋点:用户进入商品详情页、点击收藏、加入购物车、下单成功,这些事件都会被捕获并通过异步请求上报到后台。
后端接的是/api/behavior/report接口,收到事件后先校验参数,然后写入Kafka,Topic命名为user-behavior-topic。为什么不用Redis直接存而要用Kafka?因为用户行为产生的频率远高于业务请求,双11场景下每秒可能上千条写入,直接写数据库会把连接池打满。Kafka天然是削峰填谷的缓冲层,Spark Streaming从Kafka消费,批量更新用户画像。
我当时在Kafka部分配置了三个分区,副本因子设为1。伪分布式环境下副本因子设太高会导致消息积压,三副本在只有一个节点的情况下反而会让ISR同步失败。
采集数据的质量也直接影响推荐效果。前端上报的时候,要带上userId、商品ID、行为类型、行为发生的时间戳、页面来源五类信息。不带上来源的话,你无法区分用户是通过推荐位进入的还是通过搜索进入的,也就没法评估推荐位的点击转化率。
5. 可视化大屏开发实战
5.1 大屏技术选型与布局
可视化大屏是整个项目里看起来最“炫”的部分,也是我在答辩时被问得最多的一块。技术上我选了Vue 3 + ECharts 5 + DataV组件库的组合。ECharts生态成熟,社区资料多,遇到任何图表问题都能很快找到解决方案;DataV则提供了大屏常用的边框、装饰、滚动列表这些现成组件,能少写很多CSS。
大屏的布局我固定在1920×1080分辨率下开发,这也是大多数演示屏幕的标准分辨率。页面分成五大区块:顶部是整屏标题和数据时间;左上角显示宠物商品销量Top10;左中显示价格分布区间统计;中间核心区域展示地图和实时推荐热力数据;右侧是电商平台价格对比柱状图和各品牌市场占有率饼图。
每个区块我用flex弹性布局配合grid栅格定位,实测在1080p屏幕上各区块之间没有互相挤压的情况。大屏开发最忌讳的是把所有图表堆在一个页面上没有主次,用户一眼扫过去不知道看哪里。我当时把中间地图区域的面积占比设到最大的40%,让视觉焦点集中在核心数据上。
5.2 大屏数据接口与实时刷新
大屏的数据接口单独写在pet-dashboard模块里,避免和其他业务接口混在一起。接口返回的是预先聚合好的JSON结构,前端拿到后直接渲染,不在页面上做二次计算,这样可以减少前端的性能开销。
实时刷新这部分,我一开始用的是前端定时器,每10秒轮询一次后端接口。这种方式实现简单,但有两个问题:一是服务端压力大,二是数据更新有延迟。后来改成WebSocket推送,Spark Streaming每计算完一批结果就通过SpringBoot内置的WebSocket端点推送给前端,数据的实时性从10秒级别提升到2秒以内。
WebSocket端点的核心配置:
@Configuration @EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(new DashboardHandler(), "/ws/dashboard") .setAllowedOrigins("*"); } }前端通过new WebSocket("ws://localhost:8080/ws/dashboard")建立连接,收到消息后调用ECharts实例的setOption方法更新图表。这里要注意的是,每次setOption最好设置notMerge: true,否则新旧数据维度不一致时图表会出现残留的旧序列。
5.3 大屏的核心图表组件
下面是我在大屏上用的核心图表和对应的数据分析维度:
| 图表类型 | 展示内容 | 数据来源 | 更新频率 |
|---|---|---|---|
| 柱状图 | 各平台宠物商品销量排名 | MySQL销量汇总表 | 每小时 |
| 折线图 | 全网平均价格走势 | Spark聚合结果 | 每小时 |
| 饼图 | 品牌市场份额 | Hadoop ADS层 | 每日 |
| 排行榜 | 热门商品Top10 | Redis缓存 | 每10分钟 |
| 地图 | 各省份用户活跃度 | 用户登录IP解析 | 实时 |
排行榜这里我踩过一个坑:直接从Redis取全量商品列表然后在内存里做sort,第一次能跑,但数据量到10万级以后接口响应从几十毫秒涨到了快2秒。后来改成在Redis里用ZSET存储,Score就是销量,用ZREVRANGE命令直接取前10,性能瞬间回到毫秒级。
大屏的配色也是门学问。我当时选了深蓝色背景+荧光绿+橙色高亮的三色搭配,深色背景可以衬托数据的光效,荧光绿是科技感的代表,橙色用来提示关键指标,比如价格异常波动。这个配色方案在实际答辩中观众反馈不错,信息能一眼抓住注意点。
6. 调试问题排查与避坑实录
6.1 大数据组件层面的高频故障
做这个项目期间,我记录的故障排查笔记大概写了四十多条,这里挑几个最有共性的分享。
第一个是HDFS启动后DataNode进程反复退出,日志里报Incompatible clusterIDs错误。这个问题的根源是格式化NameNode之后,DataNode的storage目录里保留了旧集群ID。解决办法很粗暴:把DataNode和NameNode的目录都删掉,重新格式化,然后重启集群。这个坑在开发环境很常见,因为你会频繁切换Hadoop版本或者修改目录配置。
第二个是Spark任务日志里出现java.lang.OutOfMemoryError: Java heap space。排查思路是先确认是Executor内存不足还是Driver内存不足。如果报错出现在collect操作后,多半是Driver在收集结果时内存爆了——不建议把全量RDDcollect回Driver,应该用foreachPartition落盘或只取抽样数据。
第三个是Kafka消费者组offset提交失败。我当时做了一个延迟重试机制:
consumer.commitAsync((offsets, exception) -> { if (exception != null) { log.error("offset提交失败,等待下次提交", exception); } });这里只记录日志而不重试,是故意的。因为Kafka的commitAsync回调本身不能保证顺序,如果贸然做同步重提交,容易把后面的offset也搞乱。正确做法是允许异步提交失败,依赖下一次提交时带上最新offset,让旧的重复消费由下游幂等逻辑兜底。
6.2 SpringBoot集成层的常见问题
SpringBoot这层也和Spark/Hadoop有不少交互点,最容易翻车的是SparkSession在Spring容器里的生命周期管理。我最初尝试用@Bean方式创建SparkSession,结果每次请求都触发SparkContext already exists异常。原因在于SparkContext是JVM级别的单例,Spring的Bean每次上下文刷新都会尝试新建。
解决办法是加一个全局判断:
@Bean public SparkSession sparkSession() { SparkSession session = SparkSession.getActiveSession().isDefined() ? SparkSession.getActiveSession().get() : null; if (session == null) { return SparkSession.builder() .appName("PetRecommend") .master("spark://localhost:7077") .getOrCreate(); } return session; }另一个问题来自MySQL和Hadoop的时区差异。Spark写回MySQL的数据,时间字段经常差了8个小时。排查后确认是MySQL的serverTimezone是UTC,而Spark默认取系统本地时区。在JDBC连接串后面加上serverTimezone=Asia/Shanghai,问题就解决了。这种小问题极其隐蔽,建议所有人在项目一开始就统一时区配置。
6.3 常见问题排查速查表
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
| DataNode启动即退出 | NameNode与DataNode的clusterID不一致 | 清空数据目录重新格式化 |
| Spark任务OOM | Executor内存配置过大或Driver收集全量数据 | 调小Executor内存,避免collect |
| Kafka消息积压 | 消费者处理能力不足 | 增加分区数或消费者线程数 |
| 推荐列表为空 | 用户无行为数据且冷启动策略未配置 | 增加热门商品兜底逻辑 |
| 大屏图表数据重叠 | setOption未设置notMerge | 设置notMerge: true |
| 接口响应慢 | 数据未走缓存直接查HDFS | 增加Redis缓存层 |
| SparkSession冲突 | Spring Bean重复创建 | 使用getOrCreate单例模式 |
这张表我打印出来贴在显示器旁边,排查问题的时候按图索骥,效率提高了不少。
6.4 性能调优的实操心得
性能调优是我花时间最多的地方。最初离线推荐任务处理100万条行为数据要跑25分钟,经过以下三轮优化,压到了不到6分钟:
第一轮,数据格式从CSV换成Parquet,读取耗时降了45%。
第二轮,把频繁使用的DataFrame注册成临时视图后,多处逻辑复用spark.sql,避免了反复读源数据。
第三轮,在ALS训练前过滤掉极端用户——比如单日行为数超过500条的疑似爬虫用户,以及只有1条行为的新用户,这两个过滤条件大幅减少了训练数据量,但推荐质量并没有明显下降。
关于Spark调优还有一个很容易忽略的点:spark.sql.shuffle.partitions默认是200,这个值对中等规模数据其实偏大。我把它从200调成50,shuffle阶段的小文件数量大幅减少,整个Job的耗时降了将近三分之一。合理设置并行度比盲目加内存更管用。
7. 项目复盘与后续扩展想法
7.1 整套系统跑通后的体会
项目收尾的时候,我重新梳理了一遍整套系统的技术链路:从前端埋点采集用户行为,到Kafka缓冲,再Spark做分布式计算,HDFS做海量存储,SpringBoot做接口编排,最后ECharts渲染数据大屏。整套链路跑通用时大约四个月,其中环境搭建和排错占了将近一半的时间。
我个人的切身体会是:真正学到的不是某个单独的API怎么调用,而是怎么建立“数据从哪来、存到哪、怎么算、如何用”的全局视角。在学校里做CRUD项目时,你只管界面和增删改查,根本体会不到数据量大到一定程度后,从数据库查询到文件存储、从单机计算到分布式计算那些思维方式的转变。
比价系统的核心难点从来不是代码量,而是数据的可靠性。爬虫抓下来的数据必然有坏数据,有重复数据,有价格单位不统一的数据,清洗这一步做得不到位,后面任何GA分析、推荐算法、价格聚合都建立在沙子上。数据质量不是锦上添花,是地基。
7.2 这套系统的边界与局限
说实话,这套系统离真正商业化的比价平台还有很大距离。最大短板是数据源的覆盖度:我采集了三个平台的数据,而且每个平台只抓了大约四万个商品,而真实的宠物电商市场SKU数量远超这个量级。再一个,各平台的加密参数会频繁变化,爬虫被封也是常事,需要不断维护。
推荐效果评测方面,我只做了离线的A/B测试对比,用精确率、召回率、覆盖率三个指标衡量。在线业务的点击率提升幅度其实没有做过严谨的AB实验,毕竟环境有限。如果有条件,后续可以接入一个正规的推荐评测框架做更客观的评估。
7.3 还能往哪些方向扩展
如果让我把项目再往前推一步,有三个方向值得探索。一是引入实时计算框架Flink,替代Spark Streaming做流处理,延迟可以从秒级降到毫秒级,用户行为稍纵即逝,毫秒级的响应可以显著提升推荐新鲜度;二是用ClickHouse替代MySQL存储聚合结果和分析数据,列式存储对大屏查询的提速非常明显;三是在推荐算法里引入图神经网络或者深度因子分解机,捕获用户和商品之间更高阶的交互关系。
最后分享一个小技巧:整个项目维护的时候,建议先写好一套自动化部署脚本,把Hadoop、Spark、SpringBoot的启动统一封装成shell脚本。我后期几乎每天都要重启环境,没有脚本的话,每次手动输五六条命令,出错概率太高。自动化脚本投入一天时间,后续能省下大量重复劳动。这就是我做完这个项目最实在的经验——系统本身很重要,让系统能跑起来的流程管理,同样重要。