简介:一套基于Hadoop的协同过滤视频推荐系统项目资源,面向大数据专业的课程设计、毕业设计以及推荐系统入门开发者,重点解决海量视频数据下用户行为分析、相似度计算与个性化推荐结果的展示问题。压缩包采用zip格式,共包含383个文件,大小约12.1MB;其中58个Java源文件实现数据清洗、用户相似度计算与MapReduce统计等核心逻辑,54个CSS文件与88个JavaScript文件构成前端展示与后台管理界面,另含HTML页面、SQL脚本、图片及Spring Boot相关配置,项目的代码分层和目录结构清晰,便于导入IDE后直接阅读和二次开发。资源已吸引24人浏览学习。相比单纯的算法教程,这份资源更贴近可运行项目形态,从HDFS数据存储、离线计算到前端联动均有体现,既能作为Hadoop生态应用的综合范例,也能为协同过滤推荐系统的工程化落地提供直观参考。
1. 基于Hadoop的协同过滤视频推荐系统:这套zip到底是课程设计还是生产雏形
很多做大数据课程设计、或者准备把推荐项目写进简历的人,会在这个标题前停下来:Hadoop、协同过滤、视频推荐,三个词个个都熟,放到一起就成了黑匣子。我拿到这类zip的第一反应不是急着解压,而是先问三个问题——数据存在哪、算法跑在哪、结果写给谁。这套方案的本质,是用HDFS存用户对视频的历史行为,用MapReduce把协同过滤的矩阵运算拆到多台机器上,最后输出一份「用户可能还想看什么」的片单。它最适合两类人:一类是要交hadoop课程设计、需要完整链路能讲清楚的学生;另一类是公司里想快速验证离线ItemCF效果、但还没上Spark的工程师。至于它离线上生产还有多远,看完下面的数据组织和作业调优,你自己会有结论。
2. 先立数据地基:HDFS目录规划与ItemCF选型,为什么视频场景不选UserCF
2.1 用户行为数据怎么建表:Hive表结构与HDFS目录规划
推荐系统第一步不是写算法,而是把原始日志收拾成一张 MapReduce 能稳定读的表。视频平台的原始行为日志通常长这样:用户ID、视频ID、行为类型(play、favorite、comment、share)、播放时长占比、时间戳。生产环境一般用 Flume 或 Kafka Connect 落 HDFS,课程设计阶段拿 awk 把访问日志切成 TSV 也够用,重点是把字段顺序固定下来,避免后面每个 Mapper 都重新猜列。
我一般会在 Hive 里先建一张外部表,让表结构和 HDFS 上的日志目录解耦。这样即使 Job 失败要重跑,也不会误删原始数据。
CREATE EXTERNAL TABLE IF NOT EXISTS video_user_action ( user_id STRING, video_id STRING, action STRING, watch_ratio DOUBLE, ts BIGINT ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t' STORED AS TEXTFILE LOCATION 'hdfs:///data/video/action';这里几个选择是刻意的。用外部表而不是内部表,是因为日志会持续追加,外部表删表不影响 HDFS 文件,误操作还有后悔药。按dt分区而不是直接全表扫描,是因为推荐任务通常只回刷最近 30 天到 90 天的行为,分区裁剪能让每次作业少读一半以上的数据。TEXTFILE 格式虽然压缩比不高,但胜在能被 Hadoop Streaming 的 Python 脚本直接逐行读,对课程设计和原型验证最省事。
对应的 HDFS 目录建议这样规划,避免输出路径和输入路径混在一起:
hdfs dfs -mkdir -p /data/video/action hdfs dfs -mkdir -p /data/video/tmp/rating hdfs dfs -mkdir -p /data/video/tmp/similarity hdfs dfs -mkdir -p /data/video/output/rec输入、中间结果、最终输出分三个目录层级放,是我踩过坑之后养成的习惯。中间结果单独放 tmp,一方面是方便调试时单独重跑某一步,另一方面是防止输出目录被下一次作业的-output参数误删。很多新手把 rating 和 similarity 写进同一个父目录,跑第二次作业时直接把第一次的结果清掉了,这种翻车非常不值。
2.2 算法选型:ItemCF与UserCF,视频推荐为什么站前者
协同过滤分两类,基于用户的 UserCF 和基于物品的 ItemCF,标题没有指定用哪种,但视频推荐场景下我基本只会选 ItemCF,原因不是 UserCF 不能跑,而是它在视频站上性价比太低。
UserCF 的核心是「和你口味相似的人喜欢什么,就推给你什么」。它需要实时或近实时地维护用户与用户之间的相似度矩阵,用户量是千万级时,两两比较的计算量是用户数的平方,而且新行为一来,用户相似度就要局部重算,运维成本很高。UserCF 更适合新闻类场景,因为新闻热点变化快,用户此刻的兴趣和「相似的人」高度相关,推荐理由也往往是「和你相似的人都在看」。
ItemCF 的核心是「你喜欢的视频和哪些视频像,就把像的推给你」。它离线计算的是物品与物品的相似度,视频数量虽然也大,但相比用户量通常低一个数量级,而且视频的属性相对稳定,今天算出来的相似度明天还能用。更重要的是 ItemCF 的推荐理由天然可解释:「因为你看了《三体》第一集,所以推荐《三体》第二集」。视频站的产品经理最喜欢这种能说清楚来源的推荐。
| 对比维度 | ItemCF(基于物品) | UserCF(基于用户) |
|---|---|---|
| 离线计算量 | 与视频数相关,可控 | 与用户数平方相关,增长快 |
| 时效性 | 物品相似度更新频率低 | 用户兴趣变化快,需频繁更新 |
| 可解释性 | 强,按视频关联解释 | 弱,依赖「相似用户」概念 |
| 冷启动表现 | 新视频无人评分时差 | 新用户行为少时差 |
| 适合场景 | 视频、电商、图书 | 新闻、社区、热点内容 |
如果你的核心诉求是「让推荐理由能被答辩老师听懂」,ItemCF 也更占便宜。它只需讲清楚「共现矩阵 + 余弦相似度 + Top-N」三条线,而 UserCF 还要额外解释用户相似度矩阵的稀疏性处理,容易把自己绕晕。
2.3 Hadoop运行环境:伪分布式够用吗,什么时候需要集群
很多人一上来就搭三台机器,结果光配 SSH 和免密就耗掉一个周末。其实做课程设计或原型验证,hadoop伪分布式搭建完全够用,单机伪分布式跑通 ItemCF 全流程,再迁移到集群只是改路径和资源参数的事。
伪分布式只需要一台 4GB 以上内存的机器。核心是把 YARN 的资源参数调小,否则默认配置会试图为每个容器申请 1GB 以上的内存,小机器直接 OOM。我常用的配置是:
<property> <name>yarn.nodemanager.resource.memory-mb</name> <value>3072</value> </property> <property> <name>mapreduce.map.memory.mb</name> <value>1024</value> </property> <property> <name>mapreduce.reduce.memory.mb</name> <value>1024</value> </property>yarn.nodemanager.resource.memory-mb是 NodeManager 能支配的总内存,伪分布式下它和物理内存共享,必须留出系统余量。mapreduce.map.memory.mb是单个 Map 容器的内存上限,数据量不大时 1GB 足够,调太大会挤压同机运行的 DataNode 和 NameNode。
什么时候需要上真集群?我个人的判断标准是:单机跑一次全量 ItemCF 的时间超过两小时,或者训练数据达到亿级行为记录,再考虑用三节点 hadoop集群搭建。集群模式和多了一个必须注意的点:NameNode 要单独部署,ResourceManager 也不要和 DataNode 挤在一起;如果做 Hadoop HA,还要先搭 ZooKeeper,hadoop和zookeeper整合实战的顺序是「先 ZK 后 HDFS 再 YARN」,别反着来。用 docker 镜像做集群能省掉环境安装时间,但容器一旦重启,NameNode 的元数据目录如果没挂载宿主机磁盘,格式化信息会丢,等于白跑。
3. 核心算法拆解:评分矩阵、余弦相似度与Top-N的MapReduce实现
3.1 从原始行为日志到评分矩阵:第一个MapReduce作业
ItemCF 的输入是评分矩阵,但原始日志不是评分。用户看了视频 10 秒和看了 30 分钟,对你的推荐信号强度完全不同。我惯用的做法是在第一个作业里把行为映射成带权重的评分,把无意义的行为直接过滤掉,这样后续相似度计算会干净很多。
这一步用 Hive 做比写 MapReduce 更快,因为它是纯行转列逻辑:
INSERT OVERWRITE TABLE video_user_rating SELECT user_id, video_id, CASE WHEN action = 'play' AND watch_ratio > 0.5 THEN 1 WHEN action = 'favorite' THEN 2 WHEN action = 'comment' THEN 1 ELSE 0 END AS score FROM video_user_action WHERE dt >= '2024-01-01' AND dt <= '2024-03-31' AND action IN ('play', 'favorite', 'comment');watch_ratio > 0.5这个阈值是关键中的关键。只看过 5% 就关掉的播放,多半是误点或者内容不合口味,把它当成正样本会严重拉低推荐质量。我给视频行为赋分时,把「收藏」设为 2 分,因为主动收藏比播放完成更能表达兴趣;「完整播放」和「评论」各 1 分,都是有效正反馈。注意这里没有减分项,因为视频场景里「不喜欢」的定义很模糊,跳过不代表讨厌,贸然记负分会让相似度矩阵充满噪声。
如果不想依赖 Hive,也可以写一个简单的 MapReduce 做同样的事。Map 端解析 tab 分隔的字段,只输出(user_id, video_id)且 score 大于 0 的记录;Reduce 端不需要做什么,直接把 Map 输出透传到下一个 Job。用 Hive 的收益是省掉一个 MR 作业的提交和调度时间,缺点是让链路里多了一个「黑盒」,答辩时容易被追问 Hive 底层怎么执行。
3.2 物品相似度计算:余弦相似度的分布式聚合
评分矩阵准备好之后,下一步是算视频之间的相似度。余弦相似度的原始公式是sim(i,j) = 共同评分向量内积 / (向量i模长 * 向量j模长),在一个节点上算很简单,难的是把这个式子拆到多台机器上。
我的做法是分成两个阶段。第一阶段把「用户-视频」的评分表转成「视频-用户」倒排表,也就是同一个视频下挂上所有给它评过分的用户;第二阶段两两计算视频对。用 Hadoop Streaming 加 Python 实现是最短路径,mapper 负责输出视频对,reducer 负责累加。
#!/usr/bin/env python3 # mapper_sim.py:把评分矩阵拆成视频对 import sys for line in sys.stdin: fields = line.strip().split('\t') if len(fields) < 3: continue user_id, video_id, score = fields[0], fields[1], float(fields[2]) if score <= 0: continue # 按用户聚合后,同一个用户看过的视频两两组合 print(f"{user_id}\t{video_id}:{score}")上面的 mapper 输出还是按用户组织的,真正做视频两两组合的逻辑通常放在 reducer 里:同一个用户 ID 的所有视频到达同一个 reducer 后,在内存里做笛卡尔积。这个方案对内存有要求,如果单个用户看过的视频数超过几千,笛卡尔积会撑爆 reducer 堆内存。更稳的做法是把评分矩阵广播到每个节点,然后只在 mapper 端做分块计算,但那样代码量会膨胀。
#!/usr/bin/env python3 # reducer_sim.py:聚合视频对的共现与内积 import sys from collections import defaultdict current_user = None scores = defaultdict(float) for line in sys.stdin: fields = line.strip().split('\t') if len(fields) < 2: continue user_id, item_score = fields[0], fields[1] video_id, score = item_score.split(':') score = float(score) if current_user and user_id != current_user: # 一个用户处理完,输出该用户对所有视频对的贡献 videos = list(scores.keys()) for i in range(len(videos)): for j in range(i + 1, len(videos)): vi, vj = videos[i], videos[j] print(f"{vi}\t{vj}\t{scores[vi] * scores[vj]}") scores.clear() current_user = user_id scores[video_id] = score # 处理最后一个用户 videos = list(scores.keys()) for i in range(len(videos)): for j in range(i + 1, len(videos)): vi, vj = videos[i], videos[j] print(f"{vi}\t{vj}\t{scores[vi] * scores[vj]}")这段代码输出的每一项是「视频对的内积贡献」。下一步还有一个极简的聚合作业,对相同(vi, vj)的贡献值求和,再除以两个视频的模长,就得到余弦相似度。模长可以在评分矩阵阶段单独统计,也可以在这个 reducer 里顺带维护每个视频的平方和。
这里有一个重要的参数认知:上述 reducer 是按 user_id 分区的,mapreduce.job.reduces的值决定分区数。我一般设为 4 到 8,目的不是并行加速,而是避免单个 reducer 文件过大导致下游读取卡顿。伪分布式环境下 reduces 设太大反而会因为容器排队变慢。
3.3 生成Top-N推荐列表:排序、截断与已看过滤
相似度矩阵算出来之后,最后一步是为每个用户生成推荐列表。对某个用户而言,把他评分过的视频记为「种子集」,把种子集中每个视频的 Top-K 相似视频拉出来,按加权分排序,去掉他已经看过的,取前 N 个,就是推荐结果。
这个阶段我通常会写一个独立的 reducer,输入是相似度矩阵 + 用户种子集,输出是用户ID \t 推荐视频ID \t 得分。
#!/usr/bin/env python3 # reducer_top.py:生成每个用户的Top-N推荐 import sys from heapq import nlargest top_k = 10 # 每个种子视频取相似度最高的K个 for line in sys.stdin: fields = line.strip().split('\t') if len(fields) < 3: continue user_id, video_id, score = fields[0], fields[1], float(fields[2]) if video_id in seen_videos: # 过滤已看过的视频 continue score = score * seed_weight # 叠加种子视频权重 print(f"{user_id}\t{video_id}\t{score}")实际实现时,会把种子集作为 side data 分发给每个 mapper,然后用heapq.nlargest在每个 reducer 内维护一个大根堆,避免全量排序。这里容易被忽视的是「已看过滤」必须放在推荐输出之前,否则用户打开首页第一屏都是自己刚看完的片子,产品体验非常差。
mapreduce.input.fileinputformat.split.minsize这个参数在这里开始有意义。如果输入是大量小文件,默认的 InputSplit 会把每个文件单独分给一个 Map 任务,产生成千上万个 Map,调度开销远大于计算本身。把split.minsize调到 64MB 或 128MB,可以强制把小文件合并为较大的分片。这也是面试里常问的 inputsplit 概念:一个分片对应一个 Map 任务,分片大小直接影响 Map 并行度和资源占用。
4. 把zip落地成可运行工程:目录结构、依赖jar包与作业提交实战
4.1 解压后的工程目录:conf/src/lib/data各司其职
拿到一个「基于Hadoop的协同过滤视频推荐系统.zip」,解压后第一件事不是看代码,而是看目录结构。正常的交付工程会按照「配置、源码、依赖、数据、脚本」分层,一个典型的布局长这样:
video-recsys/ ├── conf/ │ ├── core-site.xml │ ├── hdfs-site.xml │ └── mapred-site.xml ├── lib/ │ └── hadoop-streaming-2.10.1.jar ├── src/ │ ├── itemcf/ │ │ ├── mapper_rating.py │ │ ├── reducer_sim.py │ │ └── reducer_top.py │ └── driver/ │ └── ItemCFDriver.java ├── data/ │ ├── sample_ratings.tsv │ └── sample_video_meta.tsv ├── bin/ │ └── run_all.sh └── README.md我判断一个工程能不能跑,先看bin/run_all.sh和conf/里的配置。conf里如果带着hadoop-env.sh的HADOOP_HOME变量,说明作者是让你把同一套配置拷到集群上用的;lib里有没有 hadoop-streaming 的 jar 包决定了 Streaming 作业能否直接提交。很多人拿到 zip 后嫌目录乱,直接把文件和 jar 平铺在一个文件夹,结果 ClassNotFound 满天飞,问题大多出在HADOOP_CLASSPATH没覆盖全依赖。
4.2 作业提交与资源参数:从伪分布式到集群的配置差异
Streaming 作业的提交命令看起来就一行,实际参数取舍决定成败。我建议所有路径和参数不要硬编码,全部通过 shell 变量传,这样从伪分布式切到集群只需要改顶部三行。
#!/bin/bash # run_all.sh 伪分布式一键跑通,集群环境改下面三个变量即可 export HADOOP_HOME=/opt/hadoop export HADOOP_CLASSPATH=$HADOOP_HOME/share/hadoop/common/*:$HADOOP_HOME/share/hadoop/mapreduce/* HDFS_INPUT=/data/video/action HDFS_OUTPUT=/data/video/output/rec STREAMING_JAR=$PWD/lib/hadoop-streaming-2.10.1.jar hadoop jar $STREAMING_JAR \ -D mapreduce.job.reduces=4 \ -D mapreduce.map.memory.mb=1024 \ -D mapreduce.reduce.memory.mb=1536 \ -files $PWD/src/itemcf/mapper_rating.py,$PWD/src/itemcf/reducer_sim.py \ -input $HDFS_INPUT \ -output $HDFS_OUTPUT \ -mapper "python3 mapper_rating.py" \ -reducer "python3 reducer_sim.py"-files会把本地脚本分发到集群所有节点,脚本里不要写本地绝对路径,只读 stdin 和当前目录。-D mapreduce.job.reduces=4设的是 reducer 数,不是并发度;在伪分布式下,这个值超过可分配容器数,任务会排队而不是变快。mapreduce.map.memory.mb给太小会让大文件分片的 mapper 频繁溢写磁盘,给太大会让单机上的并行 Map 数变少,需要和yarn.nodemanager.resource.memory-mb配合着调。
如果 zip 里带的不是 Streaming 脚本而是 Java 工程,提交方式会换成hadoop jar。这时最容易翻车的是 Hadoop 已编译 jar 包的版本和集群版本不一致,比如本地用 2.10 编译,集群是 3.3,反射调用接口时直接NoSuchMethodError。我的习惯是提交前先跑一个空作业验证客户端与集群版本兼容,再上完整链路,这一步能省掉半天排错时间。
4.3 结果导出:从HDFS落回MySQL,给前台一个可查询的接口
推荐结果在 HDFS 上是一堆part-00000文件,前端不可能直接读 HDFS。常规做法是把最终结果合并后导入 MySQL,让推荐接口按用户 ID 查表。
hdfs dfs -getmerge /data/video/output/rec /tmp/rec_all.tsv mysql -h 127.0.0.1 -urecsys -p123456 recsys <<EOF LOAD DATA LOCAL INFILE '/tmp/rec_all.tsv' INTO TABLE video_recommendation FIELDS TERMINATED BY '\t' (user_id, video_id, score); EOF用-getmerge而不是-cat part-*,是因为 getmerge 会正确处理文件名排序和换行拼接,尤其在 reducer 输出多个文件时不会出现半行错位。LOAD DATA 前要确保 MySQL 目标表结构里的字段顺序和 TSV 一致,否则静默截断很难查。Windows 本地如果用的是 mysql80 zip 配置的实例,还要额外确认secure_file_priv路径是否放行/tmp目录,否则 LOAD DATA 会报文件不可读。
这一步的经验是:永远先落本地临时文件再导入数据库,不要用hadoop fs -cat直接管道给 mysql 客户端。数据量大时管道会断,断点重传的成本远高于临时文件方案。导入完成后,顺手在 MySQL 里给(user_id, score)建联合索引,推荐接口的查询时间能从秒级降到毫秒级。
5. 避坑指南:Hadoop协同过滤视频推荐系统最常见的5个翻车现场
5.1 现象:NameNode起不来,日志报元数据目录不一致
格式化 HDFS 之后重启集群,NameNode 一直处于 safemode 或直接退出,错误日志里 clusterID 对不上。原因是多次hdfs namenode -format把元数据目录的 clusterID 重置了,而 DataNode 还保留着上一轮的注册信息。解决方法是把 name 和 data 两个目录下的current/VERSION文件里的 clusterID 改成一致,或者干脆停掉集群,清空dfs.namenode.name.dir和dfs.datanode.data.dir配置的所有目录,重新格式化。没有重要数据时,后者更快。这个坑的根源是把临时目录默认放在/tmp下,系统重启后/tmp被清空,元数据全丢。
5.2 现象:相似度结果全是0或NaN,推荐列表永远空
日志和代码看起来都正常,但最终输出的相似度矩阵里大部分数值是 0 或者 NaN。原因九成出在评分矩阵上:某个视频的评分向量全是 0,模长为 0,余弦相似度的分母除零。或者评分没有归一化,两个视频只被同一个用户看过,内积和模长全由那一个用户贡献,结果算出来是 1.0,不具备泛化意义。解决思路分两层:一是生产评分矩阵时过滤掉score = 0的记录,保证进入相似度计算的向量都有有效值;二是分母上加一个极小值 epsilon 兜底,Python 里写成denominator = sqrt(len_i) * sqrt(len_j) + 1e-6,从数学上根除除零。这属于典型的「数据问题伪装成算法问题」,排查时从输入数据下手比改公式更快。
5.3 现象:Reducer收不到数据,shuffle阶段缓慢或倾斜
任务卡在 reduce 阶段长时间不结束,几个 reducer 跑得飞快,一两个 reducer 卡到超时。这是数据倾斜的典型特征。相似度计算按 user_id 分区时,头部用户的行为量可能是普通用户的几百倍,所有视频对贡献集中在同一个 reducer。解决有三个抓手:第一,在 map 端加 Combiner,把同一个用户内部先聚合一遍;第二,把分区改成按视频 ID 哈希,让贡献值更分散;第三,调大mapreduce.reduce.slowstart.completedmaps到 0.9,让 map 全部完成后再启动 reduce,避免 shuffle 期间反复拉取。注意 Combiner 的输出格式必须和 Mapper 一致,否则会静默丢数据。
5.4 现象:提交作业时报ClassNotFoundException或找不到依赖jar包
作业在客户端阶段直接抛ClassNotFoundException: org.apache.hadoop.streaming.HadoopStreaming,或者java.io.IOException: No such file for streaming jar。原因基本都是HADOOP_CLASSPATH没配全,或者 zip 里的 lib 目录缺少对应版本的 streaming jar。解决方法是显式把lib/*加进 classpath,提交命令里加-libjars参数让 YARN 把依赖广播到各节点:
hadoop jar $STREAMING_JAR \ -libjars $PWD/lib/hadoop-streaming-2.10.1.jar \ -input $HDFS_INPUT -output $HDFS_OUTPUT \ -file $PWD/src/itemcf/mapper_rating.py \ -file $PWD/src/itemcf/reducer_sim.py-file和-files的区别只在命令行风格,作用都是分发本地文件。真正要检查的是 lib 目录里有没有 hadoop-common、hadoop-mapreduce-client-core 这两个基础包,很多精简过的 zip 会把它们裁掉,导致反序列化报错。这个检查应该在解压第一天就做,而不是等作业跑了十分钟才暴露。
5.5 现象:推荐结果偏科,热门视频霸榜,新视频冷启动为零
这是推荐系统最容易被非技术角色挑战的问题:推荐列表翻来覆去都是那几部热门作品,新上线的视频一个都出不来。原因是 ItemCF 天然偏向物品的流行度,共现矩阵里热门视频与所有视频都有共现,冷门视频几乎没有机会进入 Top-N。要解决需要加一把流行度惩罚,常见做法是相似度乘一个log(1 + N / popularity)的降权因子,或者在最后排序时对热门视频做降权。冷启动中的新视频没人评过,任何协同过滤都无能为力,需要并行走一条规则推荐:新视频按分类打底,播放量达到阈值后再进入协同过滤候选池。答辩时主动说出这一点,比硬扛「为什么推荐没有新品」效果好得多。
6. 从跑通到能答辩:离线指标验证与面试考点对照
6.1 用P@K和覆盖率验证推荐质量,而不是只看演示效果
能跑通只是起点,能证明推荐有效才是这个项目的价值所在。离线评估最常用的指标是 P@K 和覆盖率。P@K 的计算不复杂:把每个用户的最近行为按时间切分,前 80% 训练,后 20% 当测试集,模型给测试集的每个用户推 K 条,命中测试集里的真实观看记录就算中。Python 里几行就能算:
def precision_at_k(recommended, held_out): hit = sum(1 for item in recommended if item in held_out) return hit / len(recommended)覆盖率则是「被推荐到的视频数 / 总视频数」,它衡量推荐系统有没有把长尾内容带出来。这两个指标一个看准确度,一个看多样性,配合使用基本上能判断这个 Hadoop 方案到底能不能立住。我自己的血泪经验是:光给答辩老师演示 UI 和命令输出远远不够,一定要把指标表事先算好,一张 P@5、P@10、覆盖率的对照表,比十页代码讲解都更有说服力。
6.2 进阶方向:什么时候值得把MapReduce换成Spark
MapReduce 版本做离线 ItemCF 有一个绕不过去的缺点:每个作业之间都把中间结果落盘到 HDFS,相似度计算这种多轮迭代的算法,磁盘 IO 占到整个耗时的一半以上。换成 Spark 的动机是让它驻留内存,计算速度通常能快一个数量级。但如果你还在课程设计阶段,我不建议一上来就换,因为 Spark 的安装、调优和调试复杂度会分散你对协同过滤本身的注意力。更务实的路线是:先用纯 MapReduce 跑通全链路,把数据分区、shuffle、内存参数这些底层感觉建立起来,然后在演示环节里放一张对比表,说明「如果换成 Spark RDD 缓存,这步从 20 分钟能压到 3 分钟」。这既能展示你的系统设计意识,又不用真的去搬集群。
6.3 hadoop面试题里与这个系统直接对应的考点
答辩和面试时,这套系统最适合用来回答三类问题。第一个是 InputSplit 和 Map 并行度的关系,你可以拿自己的作业举例:输入文件 10GB,block 128MB,默认 Map 数约 80 个,每个 Map 处理一个分片,调split.minsize可以合并小分片。第二个是数据倾斜怎么解决,直接讲 5.3 节的 Combiner 和重分区思路。第三个是「为什么离线推荐不用实时框架」,答案是视频推荐对时效性不敏感,用户看视频的习惯几小时甚至几天内不会剧变,离线计算一小时更新一次足够,而且 MapReduce 的吞吐量和稳定性在离线场景优于 Flink。把这几个点串起来,这个 zip 就不只是代码,而是一个能讲透「数据、算法、运维」三层的完整项目。
这套方案做到能跑通、指标能看、瓶颈能讲,就已经值回投入。如果未来真要去生产环境,我会先补三件事:把合并后的推荐结果加上过期时间,把相似度矩阵改成每周全量重建而不是每天重算,再把用户维度换成百亿级埋点验证一遍扩容路径。方向值得做,但每一步都要踩在真实数据上。希望帮到你。
本文还有配套的精品资源,点击获取