刚接到这套视频推荐系统项目需求的时候,团队内部其实有过一段分歧:有人主张用传统的SSM框架加一个协同过滤算法库就交差,理由是用户量撑死几万人,根本用不到大数据那套组件。但实际做了预研之后我发现,这个想法在数据量小的时候没问题,可一旦用户行为日志积累到千万级,单机内存里加载用户-物品矩阵就会直接OOM,更不要说还得跑相似度计算。最后我们敲定的方案,就是标题里这套组合:用Hadoop做分布式存储底座,Spark负责海量行为数据的清洗和推荐计算,Spring Boot对外提供推荐接口服务,再配上一套可视化大屏把数据价值直观展示出来。整条链路从数据采集、预处理、算法计算到接口服务、前端展示全部打通,做出来之后,它既是一个能演示的完整系统,也是一条可以照着复现的大数据工程路径。这篇文章,我就把这套系统的架构设计、推荐算法落地过程、服务层接口封装、大屏数据可视化,以及我在调试过程中踩过的一堆坑,完整拆解给你。
1. 从一条点击日志到推荐瀑布流:先看清整条数据链路
很多人做推荐系统项目,第一个动作就是先写算法,这是本末倒置的。算法再漂亮,前面的数据进不来、后面的服务接不上,跑起来的Demo也只是一堆没人看的计算结果。我建议所有准备做这类系统的人,第一步先把数据流画清楚。推荐系统的本质,就是把用户行为数据转化为推荐结果,中间经过采集、存储、计算、服务化四个环节,每个环节用的技术组件都不一样,搞清楚各自的边界,才知道每一层该写什么。
1.1 为什么是这个“三件套”组合
先回答一个最常被问的问题:Hadoop、Spark、Spring Boot这三者在这套系统里到底分别干什么?
打个比方,Hadoop里的HDFS就像一个巨大的仓库,所有用户行为日志、视频元数据、历史计算结果都整齐地码在里面。这个仓库的优势是单台机器放不下的数据它能放下,坏了一块硬盘数据也不丢,这是推荐系统做全量计算的基础——你不能指望把所有日志都塞进一台机器的内存。
Spark则是仓库旁边的加工车间,它的工作模式是:从仓库里取原料(原始日志)、做加工(清洗、提取特征)、生产半成品(用户偏好向量、物品相似度矩阵),然后存回仓库或者送到下游。它跟仓库的本质区别是,它做的是计算,而且计算过程可以分布到多台机器上并行执行,比单机跑Python脚本快得多——这也是为什么行为数据一多,算法必须从应用服务里挪到Spark里跑。
Spring Boot则是外面的商店柜台。用户在App或网页上打开推荐频道,请求打到Spring Boot接口上,它从Redis或者HBase里把Spark预先计算好的推荐结果取出来,包装成JSON返回给前端。这一层不负责跑大计算,只负责快速响应业务请求。
这三个角色的配合关系清晰了,后面每一步做什么都会很明确。还有一个容易犯的错——把Spark计算逻辑写成Java类直接嵌在Spring Boot项目里,通过Java调Spark的API去coordinator任务。不是说不行,但这会儿你的推荐服务性能和编排能力都会受限,Spark的任务调度、资源管理、日志隔离都发挥不出来。我们最终的做法是,Spark作业以独立Jar包的形式运行,Spring Boot只负责触发和管理作业状态,两者职责分离。
1.2 一条日志从产生到变成推荐结果的完整旅程
我在这套系统里设计的是标准五步链路,每一步都能对应上具体组件和落地的数据:
- 前端埋点采集。用户在Web或者App端产生播放、点赞、收藏、搜索、分享等行为时,前端脚本组装一条结构化日志,通过HTTP接口实时上报,服务端验签后写入消息队列。考虑到项目规模,我们用的是Kafka,单机部署也能扛住演示时的吞吐量。
- 日志落盘HDFS。Kafka消费端定时拉取日志,按天分区写入HDFS路径,比如
/data/video/log/20250615/behavior.json。HDFS在这里的核心价值是海量日志的低成本存储,以及给Spark提供分布式的数据读取能力。 - Spark离线计算。每天凌晨低峰期,Spark作业启动,读取前一天的原始日志,做数据清洗、补全用户画像、计算推荐候选集,然后把结果写回HDFS或者直接同步到Redis、HBase。
- 结果入缓存。推荐结果并不是每次请求都临时算,而是预先算好放进Redis,接口层按用户ID直接取,这样响应时间能控制在100毫秒以内。Redis的key设计成
rec:user:{userId},value是JSON数组,按推荐位排序。 - 可视化大屏展示。大屏后端定时从MySQL和Redis拉取聚合指标,比如播放总量、分类占比、热门视频榜,通过WebSocket推到前端图表实时渲染。
这条链路走通之后,你会发现整个系统的每个模块边界都很干净:采集只管收数据,存储只管放数据,计算只管出结果,服务只管给接口,大屏只管做展示。后续无论你是想换掉某一个组件(比如Kafka换RocketMQ),还是扩展新的推荐策略,改动范围都控制得住。
2. 数据底座搭建:Hadoop存储与Spark批处理的核心配置
很多新人一上来就按网上的教程搭建三节点、五节点集群,最后集群装好了,发现根本没有那么多数据要处理。对于以学习和演示为主的个性化视频推荐系统,我强烈建议先用Hadoop伪分布式模式,单台服务器上跑NameNode、DataNode和Spark Standalone。伪分布式可不是“阉割版”,它的核心机制、配置项、启动流程和真集群一致,后面数据量真大了,改几行配置就能平滑扩展成多节点。这一章我把环境搭建的关键步骤和最容易掉坑的地方捋一遍。
2.1 Hadoop伪分布式搭建的要点
环境版本是我踩的第一个坑。Hadoop和Spark对Java版本要求很严格,Hadoop 3.3.x要求Java 8或11,Spark 3.x则全面要求Java 8/11/17,所以我建议统一装JDK 8,避免后期出现莫名的类加载错误或者API不兼容。操作系统就用CentOS 7.9或Ubuntu 20.04,内存至少8G,因为同时跑HDFS、Yarn、Spark和Spring Boot,内存小了下场只有OOM。
安装步骤分五步:配置SSH免密登录;下载Hadoop压缩包并解压;配置core-site.xml、hdfs-site.xml、mapred-site.xml、yarn-site.xml四个核心文件;格式化NameNode;启动进程验证。这里最需要关注的是hdfs-site.xml里的dfs.replication参数,伪分布式只有一个DataNode,副本数必须设为1,如果沿用默认的3,你会在50070端口的Web界面上看到大量的健康状态告警,因为这个副本根本找不到另外两台机器。
格式化NameNode要用hdfs namenode -format命令,很多新手格式化完启动后,发现DataNode起不来,第一反应就是重新格式化,然后继续失败,最后发现是格式化前后namenode的clusterID不一致,DataNode校验时直接拒绝连接。解决办法是:第一次格式化后,查看dfs/name/current/VERSION和dfs/data/current/VERSION里的clusterID,如果不一致,手动把它们改成一致,再重新启动DataNode。
2.2 HDFS目录设计与数据分层策略
HDFS的目录设计直接决定你后期Spark作业的读写效率,我强烈建议提前规划好分区目录,而不是把所有数据堆在一个路径下。我们最终用的是四级目录结构:/data/video/{raw,cleaned,result}/{yyyyMMdd}/{topic},第一层按数据类型分,第二层按天分,第三层按业务主题分。
/data/video/raw存前端上报的原始日志,JSON格式,保留最完整的原始字段,供排查问题用;/data/video/cleaned存清洗后的结构化行为数据,Parquet格式,列式存储,Spark读取时只扫描需要的列,I/O开销小很多;/data/video/result存推荐算法产出的结果,比如物品相似度矩阵、用户推荐列表。这个分层的意义,说白了就是让每一环节的输出都有明确的落盘位置,而不是让下一步作业去全量扫原始日志重算一遍。
顺带提一个细节:写数据的时候尽量用Spark的partitionBy按日期和话题分区写,配合Parquet格式,后续按天做增量计算时,Spark SQL可以通过分区裁剪,只读取需要的那一天数据,而不是把整个目录的数据都拉出来过滤一遍。这个优化在数据量大的时候,效果是分钟级的差距。
2.3 Spark作业的初始化与内存调参基线
Spark作业写起来不难,难的是让它稳定高效跑完。我们用的是Spark SQL加DataFrame API,好处是代码可读性好,而且Spark Catalyst优化器会自动做谓词下推、列裁剪这些优化。作业启动时最核心的参数是分配多少Executor、每个Executor多大内存、shuffle分区数。
用Spark运行一个清洗加计算任务,我的基线参考配置是这样的:
spark-submit \ --master spark://node01:7077 \ --executor-memory 2g \ --driver-memory 1g \ --executor-cores 2 \ --total-executor-cores 4 \ --class com.video.recommend.OfflineRecommend \ video-recommend.jar \ --date 20250615注意--executor-memory不要贪多,伪分布式下一台机器的物理内存就那么多,给Spark分太多,Yarn和HDFS会跟着吃紧,最后大家一起OOM。实际调试中我遇到过的一个很典型的场景:某个Spark任务在计算所有视频两两相似度时,shuffle的中间结果是原数据的十几倍,默认的spark.sql.shuffle.partitions=200每个分区积压大量数据,执行器频繁Full GC。把分区数调到800之后,任务稳定跑完,耗时也降了近一半。
还有一个容易踩的坑是Spark读写HDFS时的编码。默认的spark.sql.session.timeZone是系统时区,如果HDFS路径按北京时间命名,而Spark跑在UTC时区,读目录时日期对不上,作业会直接空跑。我会在SparkSession启动时显式设置spark.sql.session.timeZone=Asia/Shanghai,避免这种潜在问题。
3. 推荐算法工程化落地:ItemCF从公式到Spark代码
推荐算法是这套系统的灵魂,也是面试官一定深挖的模块。我不打算把协同过滤的公式从教科书上搬一遍,而是讲讲它怎么在Spark里真正跑起来,以及为什么在这个场景下我们最终选了物品协同过滤(ItemCF),而不是用户协同过滤(UserCF)或者ALS隐语义模型。
3.1 视频场景下为什么是ItemCF胜出
推荐系统的两大经典流派是UserCF和ItemCF:UserCF是找和你偏好相似的用户,把他们喜欢的东西推荐给你,适合新闻资讯这种用户兴趣变化快的场景;ItemCF是找和你历史喜欢物品相似的物品,核心是计算出“和这个视频相似的其他视频”,适合视频、电商这类物品数量相对稳定、用户行为能充分体现物品关联的场景。
视频推荐为什么Pick ItemCF?最大理由是视频的数量和用户数量相比要少好几个量级。一个视频平台,物品(视频)可能就十万级,用户却有百万级。ItemCF只需要维护一个物品-物品相似度矩阵,矩阵的大小是物品数的平方,十万物品算出来是一百亿个元素,Spark跑起来压力尚可;而UserCF要维护用户-用户矩阵,百万用户就是万亿级元素,单机根本不用想。另一个理由是视频的受欢迎程度波动极大——爆款的相似物品往往也是爆款,物品之间的相似关系相对稳定,可以提前离线算好,不必实时更新。
3.2 用户行为打分:播放时长和完播率如何转成数字
协同过滤的第一步是构建“用户-物品”评分矩阵,但视频平台几乎不会让用户给视频打分,我们手里能用的只有用户行为日志。这时候需要设计一套行为分数权重规则,我参照行业常用逻辑做了一张映射表:
| 行为类型 | 权重 | 说明 |
|---|---|---|
| 播放行为 | 0.5 + 完播率系数 | 播放权重=0.5,播放时长达到总时长90%以上时加0.5,达到50%以上加0.2 |
| 点赞 | 1.0 | 明确的正反馈信号 |
| 收藏 | 1.5 | 比点赞更强的偏好信号,说明用户想以后再看 |
| 分享 | 2.0 | 强正反馈,通常意味着内容质量高 |
| 评论 | 1.2 | 参与互动,但评论也有负面内容,权重要控制 |
| 不感兴趣/滑走 | -1.5 | 负反馈,降低该视频的推荐优先级 |
具体到Spark代码里,我会从清洗后的行为表里,按用户和视频分组,把同一天内同一用户对同一视频的所有行为得分加权累加,形成一张用户-视频评分明细表。这里有个容易忽略的细节:播放权重要和商品购买这类行为做区分,购买是强正反馈(权重可以给到5以上),但播放只能算弱正反馈,因为用户可能只是误点进来看一眼就退出。所以在聚合播放行为时,我会额外校验播放时长字段,播放低于10秒的,直接视为无效行为过滤掉,否则它会影响整个评分矩阵的质量。
3.3 Spark算子操作:共现矩阵与相似度TopN的完整逻辑
ItemCF在Spark里落地,核心三步是:算物品共现矩阵、归一化得到相似度、结合用户历史生成推荐列表。以下是第一步和第二步的关键伪代码逻辑,你完全可以照着实现:
// Step 1: 读取用户-物品评分表,过滤无效行为 val df = spark.sql("SELECT user_id, video_id, score FROM behavior_score WHERE score > 0") // Step 2: 按用户分组,对同一用户看过的所有视频做笛卡尔积,统计共现次数 val cooccurRDD = df.rdd .map(row => (row.getInt(0), row.getInt(1))) // (userId, videoId) .groupByKey() .flatMap { case (_, videoIds) => val distinct = videoIds.toSet.toList for { i <- distinct j <- distinct if i < j && i != j } yield ((i, j), 1) } .reduceByKey(_ + _) // Step 3: 计算每个视频被偏好过的用户数,用于归一化 val videoUserCount = df.rdd .map(row => (row.getInt(1), 1)) .reduceByKey(_ + _) val userCountMap = videoUserCount.collectAsMap() // 广播变量优化 // Step 4: 计算余弦相似度 val simRDD = cooccurRDD.map { case ((a, b), cnt) => val sim = cnt / math.sqrt(userCountMap.getOrElse(a, 1).toDouble * userCountMap.getOrElse(b, 1).toDouble) (a, (b, sim)) }这里最核心的优化是userCountMap的广播变量——所有Executor节点都只需要这个只读Map的一份拷贝,而不是每个Task都从Driver拉取,能省下大量网络传输开销。还有一个工程细节:共现矩阵里会出现某些超级热门视频和其他视频的共现次数特别高,直接用它做相似度会导致推荐结果永远集中在热门内容上,实际效果就是“每个人看到的都是同一个爆款合集”,个性化完全丧失。
我的处理方式是在相似度计算时对热门物品做降权处理,俗称alpha惩罚项:sim = cnt / pow(count_i, alpha) / pow(count_j, 1-alpha),alpha取0.5到0.8之间,具体值可以通过看推荐结果多样性调优。这个细节,你在绝大多数教程里看不到,但它是决定推荐效果是“千人一面”还是“千人千面”的关键一刀。
3.4 冷启动问题的兜底方案
任何协同过滤系统都绕不开冷启动:新注册用户没有任何行为记录,相似度矩阵对他来说是一张白纸;新上架的视频一开始也没有用户互动,不会被任何推荐链路捞出来。我们做了三层兜底:
第一层是用户冷启动处理。新用户或行为记录少于3条的用户,直接返回热门视频榜前50,而且要求热门榜必须混合多个视频分类,避免新用户第一屏全是同一个品类的内容——那会极大降低留存。第二层是新物品冷启动处理。新视频发布后,先根据它的标题、简介、分类标签做内容特征抽取,用TF-IDF向量和已有视频算内容相似度,找到和它内容最接近的视频,挂在相似推荐列表里,让它有初始曝光机会。第三层是行为充足后的个性化覆盖。当用户有效行为超过5条,系统切到ItemCF个性化列表。
这三层策略用一张开关表配置在Redis里,比如cold_start:threshold=5这样的配置项,运营同学可以随时调整阈值,不需要改代码重新发版。
4. Spring Boot服务层:把Spark算好的结果封装成可用接口
算法层产出的是离线计算结果,用户打开App瞬间请求的是一个毫秒级的接口,这中间的接缝就是Spring Boot服务层。我见过不少项目把推荐逻辑全部塞进一个controller里,里面直接调Redis、连MySQL、算集合,最后controller几百行,维护起来十分痛苦。正确的做法是按数据流向拆成接入层、业务层、数据访问层,每一层的职责严格隔离。
4.1 后端模块划分与Redis缓存策略
我们把后端拆成了四个模块:用户服务、视频服务、推荐服务、大屏服务。推荐服务是最核心的,它对外提供两个方法:一是getRecommendItems(userId, page, size)从Redis中读推荐候选列表,做分页和排序权重处理;二是getSimilarItems(videoId, limit)返回当前视频的相似推荐。两个方法都挂在这套统一返回结构上:
{ "code": 200, "message": "success", "data": { "items": [ { "videoId": 100234, "title": "Spark入门到大神实战", "coverUrl": "https://... ", "score": 0.873, "reason": "与你收藏的《Hadoop核心原理》相似" } ] } }Redis缓存是这套接口性能的关键。我把推荐结果分成热榜和个性化两类,key分别是rec:hot:list和rec:user:{userId},个性化缓存设置4小时过期,因为离线计算每天跑一次,缓存过期后下一次请求会触发一次“缓存穿透检查”——先查Redis有没有,没有就查MySQL里存放的结果表,再回填Redis。这里有个必须注意的细节:缓存击穿。如果某个热门用户的key刚好在过期瞬间有大量并发请求,所有请求都会打到MySQL上。解决方式是加一个互斥锁或者直接对热点key延长过期时间,我选择的是后者,简单有效。
4.2 定时调度:Spring Boot如何触发Spark作业
这里有一个设计取舍问题:Spark作业是独立的Jar包,那Spring Boot怎么调度它呢?我们最终用的是Spring的@Scheduled定时任务框架,每天凌晨2点,通过调用Shell脚本形式执行spark-submit命令,作业状态写入MySQL的一张recommend_job_log表。启动前插入一条RUNNING状态记录,作业结束回调更新为SUCCESS或者FAILED,然后给运维人员返回一条站内通知。
这个设计最大的隐藏坑是Shell脚本的执行超时。默认的Process.waitFor()会一直阻塞等待子进程结束,一旦Spark任务卡死,定时任务的线程就永远占着不放,第二天凌晨任务再次触发时,上一个线程还活着,两个Spark任务同时抢资源,直接把Yarn打挂。我的解法是给Process.waitFor(long, TimeUnit)设置超时时间,比如12小时还跑不完,就强制destroyForcibly()杀掉进程并标记失败。
另一个值得一提的细节是,定时任务的时区设置一定要和Spark侧一致,我前面提到过Asia/Shanghai的设置,Spring的@Scheduled也要显式指定zone = "Asia/Shanghai",否则夏令时切换或者部署在海外云主机上,定时任务会以服务器默认时区为准,出现过任务每天提前8小时跑的情形,数据完全对不上。
4.3 接口并发保护与会话保持
推荐接口属于高并发读接口,生产环境需要做限流和降级,但开发阶段至少要把基本保护做好。我们在接入层加了一个简单的令牌桶限流器,单台实例QPS超过500时,多余的请求直接返回兜底热门推荐,而不是让其积压在线程池里慢慢超时。这个保护能让系统在瞬间流量突增时,核心功能不至于彻底不可用。
另外建议接口层使用@Transactional的地方要慎之又慎。推荐列表读取只有纯查询,没必要开事务;反而是用户点击推荐位之后的上报接口(用于计算推荐转化率)才需要事务,因为它要更新行为明细表和推荐曝光统计表。用事务能保证这两张表要么同时成功要么同时回滚,避免统计数字对不上。
5. 可视化大屏:数据指标拆解与实时刷新实现
做这套系统时,可视化大屏的要求比较明确:既能实时展示视频平台的核心运行指标,又能让非技术人员一眼看懂推荐系统的实际效果。很多大屏做出来之所以被业务同学吐槽“看不明白”,就是因为它把不相关的指标堆在一起,只求看起来色彩丰富,没有逻辑层次。
5.1 大屏放什么指标,为什么放这些
我眼中的推荐系统大屏,核心要回答四个问题:平台整体活跃度如何、用户在看什么内容、推荐系统是否在起作用、系统有没有异常。围绕这四个问题,我选了六个指标卡片:
- 今日播放总量与环比:体现平台整体活跃度。
- 实时播放趋势曲线:近24小时每分钟播放量,直观反映高峰时段。
- 视频分类热度TOP5:当前点击量最高的几个内容分类,用于内容运营方向判断。
- 热门视频实时榜:播放量最高的前10条视频,配封面和点击量。
- 推荐转化率:推荐位点击次数 / 推荐位曝光次数,这是衡量推荐系统效果最核心的指标。
- 实时行为流:最近20条用户行为事件的滚动列表,展示用户参与到系统的实时态。
为什么不做成十几个图表?因为大屏的阅读时间是几十秒,人的注意力只能抓住六七个信息点,超过这个数量,图表之间互相干扰,信息密度太低,等于什么都没展示。如果你给其他平台做类似的大屏,务必把指标数量控制在12个以内,且指标之间要有业务逻辑关联。
5.2 大屏后端聚合接口与WebSocket推送
大屏数据来源分两类:一类是分钟级聚合数据,直接从MySQL的聚合表读;另一类是秒级实时数据,需要从Redis的计数器里实时读取。我们在大屏服务模块里做了一个统一的聚合接口/api/big-screen/overview,它一次返回所有指标卡片的数据,前端图表初始化时调一次,之后依赖WebSocket持续接收增量更新。
WebSocket推送的细节值得展开聊聊。最初我用的方案是前端每5秒轮询一次接口,逻辑简单,但存在秒级延迟,大屏上看到的行为流更新不够“实时”,演示效果比较拉胯。后来改成WebSocket长连接,后端用定时任务每3秒推送一次增量数据,有效负载很小,理论上单台服务器能同时撑几千个连接。注意在Spring Boot里搭建WebSocket服务时,需要额外配置线程池,自定义ServerEndpointExporter,并且处理好浏览器断线重连的逻辑,否则操作人员切一下页面,连接就永久失效了。
5.3 前端图表的选型与渲染优化
大屏前端我们用的是Vue加ECharts,选ECharts而不是其他图表库,原因是它对大数据量折线图、地图、热力图这些场景的处理最成熟,配置项丰富,社区案例多,遇到问题很容易查到解决办法。整个大屏的布局采用栅格化系统,头部放平台名称和实时时钟,中间核心区域放播放趋势图和分类热度图,左右两栏放排行榜和实时行为流,底部放推荐转化率——整体用一种由总到分的阅读顺序。
有一个很容易被忽视的性能坑:ECharts在频繁更新数据时,如果直接setOption整个覆盖,会引起图表重绘闪烁,尤其在大屏上非常难看。我的做法是单独维护数据序列,每次更新只调用setOption({ series: [{ data: newSeriesData }] })更新序列数据,而不是整个option对象。另外一个经验是,大屏滚动的实时行为流不要用console.log频繁打印,浏览器控制台打印多了会拖垮渲染线程,导致所有图表的动画卡顿。
6. 分布式环境下的调试实录:那些文档里查不到的坑
做这个项目前后花了差不多三周,其中有一半时间用在和环境、任务、数据较劲上。这一章我挑几个印象最深的坑,完整记录下来。它们不是算法问题,也不是技术原理问题,而是分布式环境下组合使用这几个组件时,才会暴露出来的边角问题。照着这套系统做的人,大概率也会遇到同样的场景,提前知道能省下大量排查时间。
6.1 NameNode安全模式与DataNode磁盘问题
项目做到一半,有一次我重启了整台开发服务器,然后访问HDFS Web界面,发现界面提示NameNode is in safe mode。处于安全模式下,HDFS只允许读操作,不允许写入,但Spark作业照常启动的话会一直报错,抛出FileNotFoundException。这个坑的常见原因是HDFS在异常关闭后重新启动,安全模式会开启一段时间等待DataNode汇报数据块,正常情况下几十秒自动退出。但如果你在安全模式期间就强行提交Spark作业,会看到一个莫名其妙的“目录不存在”或“文件找不到”异常,完全不会提示安全模式的问题。
解决方法是等待自动退出,如果长时间卡着,执行hdfs dfsadmin -safemode leave强制退出。另外我遇到过DataNode进程起不来的情况,排查了一圈发现是磁盘空间不足,DataNode的dfs.datanode.data.dir指定的目录满了,它就会自动拒绝启动——因为数据节点写不了新块。看日志你就明白了,报错信息里明确写着there is not enough space。所以伪
分布式环境,最好给数据目录预留至少20G,别把HDFS数据目录放在系统盘,因为系统运行日志会把空间吃光。
6.2 Spark作业OOM和数据倾斜的实战表现
Spark作业OOM和单机OOM的表现不太一样。我们第一次跑ItemCF共现矩阵时,遇到的是Executor异常退出,Web UI上显示某个Executor状态为FAILED,日志里能看到OutOfMemoryError堆栈,但问题不在Executor总内存大小,而是某个任务处理的数据量特别大。追根溯源,是数据倾斜——某个热门视频的共现对其他所有视频都产生了关联行,数据密度比其他视频高出一个量级,单个Task处理它所在的分区时直接扛不住。
数据倾斜是分布式计算里最容易被低估的问题。排查方法很简单,看Spark Web UI里各个Task处理的数据量,如果出现一个Task处理几百MB而其他Task只有几MB,基本就可以判定是倾斜了。常用的治理手段:一是提前过滤热点数据,比如用户量极少的异常视频不参与共现;二是对倾斜的key加随机盐值,把一个大key拆成多个小key分散到不同分区,最后再合并;三是控制单分区数据量,将spark.sql.shuffle.partitions适当调大。我们最终用了第二种,把热门视频的相似度计算单独拆出来一个作业跑,隔离问题,主作业的任务时长一下子就恢复正常了。
6.3 中文乱码和时区问题:越基础越容易翻车
最后记录一个特别容易被忽视的问题——中文字段乱码。Spark读取HDFS上的JSON文件时,如果你没有显式指定编码,spark.read.json的默认行为可能无法正确处理某些中文字段,尤其是在服务器locale不是zh_CN的情况下。我们处理过一次视频标题变成“????”的乱码,排查半天发现是服务器环境变量LANG=en_US.UTF-8导致的编码不一致。这是大数据栈里典型的“环境归因”问题,注意保证所有节点的LANG、Hadoop的io.compression.codecs等环境配置一致即可。推荐的做法是Spark作业统一用spark.sql.parquet.writeLegacyFormat和显式字符集,并且不要在代码里硬编码中文字符串到HDFS路径里,路径用英文加日期就够了。
时区问题前面已经提过两次,这里再补充一个重要场景:你本地开发机可能是北京时间,而云服务器默认是UTC时间,大屏展示的“今日播放总量”如果直接按服务器当前日期取数,一过UTC零点就会清零,等于每天提前8小时刷新,业务数据全乱。处理时我建议后端从网关层开始统一使用带时区的UTC加8时间戳,BigQuery或者MySQL存储用DATETIME带SYSTEM格式,前端展示时再转成本地时间,这样做能规避各种边界情况。
6.4 推荐结果评测:不能只看“效果不错”
项目收尾时,推荐效果往往面临“到底怎么证明有效”的灵魂拷问。完整的推荐系统评测要用准确率、召回率、覆盖率、多样性这些指标做离线评估,比如我们给每个用户划分了训练集和测试集,用训练期间的TopN推荐命中测试集的实际观看记录来计算离线命中率。实际操作时,因为这只是一个演示项目,用户交互的真实点击流样本不够,我们退而求其次,用了推荐转化率这一个在线指标辅助判断:推荐位曝光1000次,实际被点击了多少次。如果转化率稳定超过8%,基本说明推荐内容和用户兴趣是有正相关的。
同时我把每个推荐位的内容附带了推荐理由,比如“与你收藏的XX相似”这种文本,让用户和大屏观察者都能感知到推荐结果不是随机出的。这个小设计在演示和答辩时贡献非常大——所有人都能直观感受到个性化推荐的存在感,而不只是听你口头介绍“用了ItemCF算法”。
我个人的体会是,这类基于大数据的推荐系统项目,真正的高价值不只在算法本身,而在于把“数据采集—存储—计算—服务—展示”这条链路完整打通的工程能力。当你在调试中习惯了面对OOM、数据倾斜、时区错乱、安全模式各种组合问题,你会发现以后做任何大数据项目,心里都有底,知道问题大概出在哪一层、该用什么手段去求证。耐心点把这些坑一个个填平,你做出来的系统就真的是能跑、能看、能讲的完整作品,而不是停留在PPT上的架构图。