ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于Spark的个性化短视频推荐系统全流程实战解析

基于Spark的个性化短视频推荐系统全流程实战解析 1. 先拆清楚这个毕设到底在做什么不是“推荐算法”是推荐系统全流程先说结论。基于 Spark 的个性化短视频推荐系统是一个典型的计算机毕业设计项目但它并不等于让你从头发明一个推荐算法。它更接近一个“完整系统”的工程化练习把用户行为数据采集进来用大数据的思路做离线加工和特征统计再用推荐逻辑从海量短视频中筛选出当前用户可能感兴趣的 TopN 结果最终通过 Web 页面把结果展示出来。这套组合之所以常见是因为它把几个经典技术点都串起来了Python 负责整个系统的主逻辑包括推荐脚本、业务接口和数据处理脚本。Spark 负责对用户行为日志、视频特征表做分布式计算和离线批处理。Hadoop 提供底层的分布式文件存储能力也就是 HDFS用于存放行为日志和视频特征数据。Django 负责 Web 展示层和前后端接口把推荐结果以网页形式呈现给用户。前端部分一般使用 HTML、CSS、JavaScript配合 ECharts 做统计图表展示。所以它其实是一个“数据采集 离线计算 推荐逻辑 Web 展示”的完整闭环。这个闭环比单纯的算法实验要有价值得多也正好覆盖了答辩时需要讲清楚的重点数据从哪来、存在哪、怎么计算、结果怎么用。这个项目比较适合两类人。一类是做毕业设计、需要把大数据技术栈和 Web 开发串起来的本科生另一类是刚学完 Spark 和 Django、想找一个能放在简历上的综合项目的人。如果你只对推荐算法本身感兴趣那这个项目可能偏重工程化算法部分反而只是一个环节。2. 关键技术栈之间的关系先理清再动手避免“装了一堆环境却不知道谁在干活”不少人在这个项目上翻车不是因为功能有多难而是没搞清楚 Hadoop、Spark、Django 各自负责什么。你先在脑子里把这条链路建立起来Python 脚本生成或接收用户行为日志写入 HDFS。Spark 从 HDFS 读取日志经过清洗、统计、特征计算生成用户推荐候选集或短视频特征表。Django 后端通过读取 Spark 输出的结果表再按当前用户的 ID 取出推荐结果渲染到页面上。前端页面展示短视频列表和用户行为统计图表。换句话说Hadoop 是存储底座Spark 是计算引擎Django 是展示入口。这三者之间不直接竞争而是顺序配合。2.1 Hadoop 在项目里的真实位置不要把它理解成一个“必须启动很多进程”的负担Hadoop 在这个项目里的核心组件是 HDFS。你可以只在伪分布式模式下运行 NameNode、DataNode 这两个核心进程再加上 YARN 的 ResourceManager 和 NodeManager 配合 Spark 提交任务。如果磁盘和内存有限甚至可以只保留 HDFS 用于存数据Spark 在 local 模式下直接读取本地文件。这里先说一个常见误区很多同学以为必须搭一个完整的 Hadoop 集群才能做这个项目。其实不是。毕业设计阶段Hadoop 伪分布式模式已经完全够用。你只有一台笔记本内存 8GB 或 16GB也照样可以把整个流程跑通。真正的分布式集群是生产环境的事情不是毕设阶段必须面对的问题。Hadoop 在这个项目里要完成的核心任务只有两个提供一个目录结构清晰的存储位置比如/recommend/user_action/存放用户行为日志/recommend/video_feature/存放视频特征数据。让 Spark 能从 HDFS 上读取数据并把计算结果写回 HDFS 或直接写回 Django 可读取的目录。2.2 Spark 在项目里的真实位置离线批处理而不是实时计算这个项目的推荐链路通常是离线计算模式。用户行为日志会积累一段时间然后 Spark 在凌晨或定时任务里跑一次批量计算生成新的推荐结果表。Spark 在这里做的事包括数据清洗去掉重复点击、异常浏览时长、空用户 ID 等无效记录。特征统计统计用户的点击次数、观看时长、点赞数、评论数、视频类别偏好。相似度计算基于视频的类别、标签、描述文本计算视频间相似度。候选集生成从所有短视频里筛出用户可能感兴趣的候选内容。TopN 排序按综合评分输出每个用户最终的推荐列表。你可以用 RDD 或 DataFrame 完成这些操作。我更建议用 DataFrame 和 Spark SQL因为代码更简洁调试时也容易看清数据结构。许多报错其实出现在 RDD 转换环节DataFrame 的方式会少踩一些类型相关的坑。2.3 Django 在项目里的真实位置推荐结果的出口不是推荐算法本身Django 负责的是用户登录、推荐结果展示、测频偏好选择、行为反馈模拟和后台管理。你需要做的核心页面包括登录注册页确定当前用户是谁。首页推荐流展示当前用户的个性化推荐短视频列表。视频详情/播放页展示视频信息同时记录一次用户行为。行为数据页展示用户自己的点击、点赞、评论、收藏历史。后台管理页管理短视频数据和系统用户。Django 模型层的核心表至少要有用户表、视频表、用户行为表、推荐结果表。你不需要在 Django 里自己实现相似度算法或排序算法那部分由 Spark 算好Django 只需要读取结果。注意Django 和 Spark 的数据交互最稳妥的方案是“结果落库”。Spark 计算完成之后把推荐结果写入 MySQL或者在本地生成 CSV 文件再由 Django 模型导入。不要尝试让 Django 直接调用 Spark 任务来做同步推荐这样会拖垮页面响应时间也容易让答辩现场演示变得不可控。3. 环境准备与版本匹配这一节占掉了不少同学的“大半周”如果只是写算法一个 Jupyter Notebook 就够了。但这里是完整系统所以环境准备环节要足够仔细。我按实际操作顺序整理成一套建议。3.1 基础软件清单软件学习环境建议生产环境建议操作系统Windows 10/11 或 Ubuntu 20.04/22.04Linux 服务器JDKJDK 1.8 或 11与 Hadoop、Spark 版本匹配的 JDKHadoop3.x尽量选稳定版3.x 集群Spark3.x本地模式或 standalone 模式3.x 集群Python3.8 或 3.103.8 或 3.10Django4.x4.xMySQL5.7 或 8.08.0代码编辑器VS Code 或 PyCharm不做限制Windows 上开发时Hadoop 和 Spark 都能以本地模式运行。但你要特别注意两个问题。第一个是 Java 环境变量。Hadoop 和 Spark 都依赖 JDK如果JAVA_HOME没有配置对启动时会直接报错。建议先用命令行确认echo %JAVA_HOME% java -version第二个是 Windows 下 Hadoop 需要winutils.exe。如果你采用伪分布式或从本地文件读取数据可能需要下载对应版本的winutils.exe和hadoop.dll放到 Hadoop 安装目录的bin文件夹下。否则启动 DataNode 时可能出现权限或本地库报错。这个细节很多人卡了很久。3.2 本地模式跑通 Spark 的验证方法在配置整个项目之前我建议你先用 Spark 自带的示例任务验证环境是否可用。spark-submit --master local[2] examples/src/main/python/pi.py 10如果能在几秒内输出一个Pi is roughly开头的日志就说明 Spark 本地模式运行正常。这一步很关键因为后面所有 Spark 任务都要依赖这个基础能力。3.3 Python 虚拟环境与 Django 项目初始化推荐用conda或python -m venv创建独立虚拟环境。原因很简单Django 项目、Spark 的 PySpark 环境、数据处理脚本可能使用不同版本的依赖包混在一起容易出现依赖冲突。python -m venv venv venv\Scripts\activate pip install django pyspark pandas django-admin startproject recommend_web cd recommend_web python manage.py startapp user_center python manage.py startapp video_center python manage.py startapp recommend_center这里按照业务模块拆了三个appuser_center管用户video_center管短视频数据recommend_center管推荐结果展示。这种结构在答辩时也能体现出你对项目模块划分的思考。4. 数据从哪里来数据集生成与初始化很多同学在这个项目上遇到的第一个实际问题就是“我没有真实的短视频用户行为数据”。有几种方案可选。4.1 真实公开数据集如果优先考虑真实性可以找用户行为日志类的公开数据集。但是公开数据集通常字段格式和你自己定义的表结构不完全一致需要写数据清洗和转换脚本。如果你的论文里强调“使用真实数据处理”那么公开数据集是必须的。4.2 模拟数据生成脚本我个人的建议是先用模拟数据把整个链路跑通再决定是否引入真实数据集。模拟数据的好处是你可以精确控制用户量、视频量、行为分布和异常情况方便测试推荐结果是否合理。你可以用 Python 脚本生成import random import csv user_ids [fuser_{i} for i in range(1, 101)] video_ids [fvideo_{i} for i in range(1, 501)] categories [科技, 美食, 影视, 运动, 游戏, 音乐, 旅行] with open(user_action.csv, w, newline, encodingutf-8) as f: writer csv.writer(f) writer.writerow([user_id, video_id, action_type, watch_seconds, timestamp]) for _ in range(20000): writer.writerow([ random.choice(user_ids), random.choice(video_ids), random.choice([click, like, comment, favorite]), random.randint(0, 300), f2024-0{random.randint(1,6)}-{random.randint(1,28)} {random.randint(0,23)}:{random.randint(0,59)}:{random.randint(0,59)} ])这段脚本会生成 2 万条用户行为记录。100 个用户、500 个短视频、2 万条行为足够验证完整流程也不至于让 Spark 任务跑太久。4.3 数据入库流程模拟数据生成之后你需要做两件事把 CSV 文件上传到 HDFS 指定目录给 Spark 读取。把短视频基础信息导入 MySQL给 Django 展示。HDFS 上传命令hdfs dfs -mkdir -p /recommend/input hdfs dfs -put user_action.csv /recommend/input/ hdfs dfs -ls /recommend/input/-mkdir -p会递归创建目录-put负责上传本地文件到 HDFS-ls用于确认文件确实存在。这一步如果失败后面 Spark 任务读取数据时也会失败。5. 推荐逻辑的两种实现思路从简到难我不建议一开始就写复杂的机器学习公式。先跑通一个能解释清楚的推荐逻辑再根据时间和能力决定是否升级模型。下面给出两条路径。5.1 基于标签偏好的召回 热度排序这是最容易在答辩时讲清楚的方案。核心思路是统计每个用户对不同短视频类别的历史点击次数计算偏好权重。从短视频表里召回用户偏好类别下的视频。按视频的热度值比如播放量、点赞率、评论数综合排序筛选前 N 条。把结果写入推荐结果表。用 Spark DataFrame 可以这样表达from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, desc, row_number from pyspark.sql.window import Window spark SparkSession.builder.appName(VideoRecommend).getOrCreate() user_action spark.read.csv(hdfs://localhost:9000/recommend/input/user_action.csv, headerTrue, inferSchemaTrue) video_feature spark.read.csv(hdfs://localhost:9000/recommend/input/video_feature.csv, headerTrue, inferSchemaTrue) # 统计用户-类别偏好 user_cat_pref user_action.join(video_feature, video_id) \ .groupBy(user_id, category) \ .agg(count(video_id).alias(click_count)) # 对每个用户偏好类别排序 window Window.partitionBy(user_id).orderBy(desc(click_count)) user_top_cat user_cat_pref.withColumn(rn, row_number().over(window)) \ .filter(col(rn) 1) \ .select(user_id, category) # 候选视频热度排序 video_hot video_feature.withColumn(hot_score, col(play_count) * 0.6 col(like_count) * 0.3 col(comment_count) * 0.1) # 用户-偏好类别召回热视频 result user_top_cat.join(video_hot, category) \ .orderBy(user_id, desc(hot_score))最后把result写出为 Django 能读取的格式或者写入 MySQL。如果写入 MySQL可以使用pymysql批量插入。如果只是简单演示可以直接写 CSV。5.2 基于协同过滤的思路如果论文要求算法有更多理论性可以考虑基于物品的协同过滤。核心逻辑是根据用户行为矩阵计算物品之间的相似度。对某个用户看过的视频找出相似视频作为候选。按相似度和用户历史行为加权生成推荐列表。这个思路比标签偏好方案复杂一些但也不是很难。主要用到 Spark 的向量化和矩阵运算。要注意的是协同过滤结果是否合理高度依赖行为数据的丰富程度。如果每个用户只有几条行为记录物品相似度矩阵会非常稀疏推荐效果可能不理想。建议至少生成每个用户 50 条以上行为记录。6. Django 展示层设计不要只做“能显示数据”到了 Django 这一层项目就开始有一个完整的“作品”感了。这部分如果做得粗糙前面 Spark 写得再好也不容易体现出来。6.1 模型设计video_center的短视频表建议包含这些字段video_id视频唯一标识title标题cover_url封面图地址play_url播放地址category分类play_count播放量like_count点赞量comment_count评论量user_center的用户表可以保留 Django 默认的User再扩展一个字段比如昵称。recommend_center的推荐结果表则可以这样设计class RecommendResult(models.Model): user models.ForeignKey(User, on_deletemodels.CASCADE) video models.ForeignKey(Video, on_deletemodels.CASCADE) score models.FloatField(default0.0) reason models.CharField(max_length255, blankTrue) create_time models.DateTimeField(auto_now_addTrue) class Meta: ordering [-score]reason字段比较实用。比如 Spark 计算后可以写入“根据你对科技类视频的偏好推荐”这样前端可以展示推荐解释答辩时也更有说服力。6.2 推荐结果落库脚本Spark 算完后推荐结果写 MySQL 的脚本可以放在 Django 的management/commands里也可以独立运行。我常用的一种方式是把 Spark 结果输出成 JSON 或 CSV再用 Django 的loaddata或自定义脚本导入。另外要注意Django 默认的sqlite3数据库在数据量小的时候没问题但如果你用 Spark 写入大量推荐结果建议直接配 MySQL。Django 数据库配置修改很简单DATABASES { default: { ENGINE: django.db.backends.mysql, NAME: recommend, USER: root, PASSWORD: your_password, HOST: 127.0.0.1, PORT: 3306, } }然后执行迁移python manage.py makemigrations python manage.py migrate6.3 前端页面与接口Django 推荐结果页面的核心逻辑很简单根据当前登录用户的 ID从推荐结果表里查出视频列表再交给模板渲染。视图函数示例from django.shortcuts import render from recommend_center.models import RecommendResult def index(request): if not request.user.is_authenticated: return redirect(/login/) rec_list RecommendResult.objects.filter(userrequest.user).select_related(video)[:20] return render(request, recommend/index.html, {rec_list: rec_list})页面模板里用for循环渲染视频卡片即可。如果要做 ECharts 图表展示用户行为分布可以写一个接口返回 JSON再由前端图表组件读取。比如统计当前用户不同类别的点击次数分布前端展示成饼图或柱状图。这一部分在答辩演示时非常加分。7. 批量任务与调度不要让计算和展示“绑死”在同一个进程里推荐结果计算不会每刷新一次页面就重新跑一遍。那样既不现实也会让页面响应变得非常慢。正确做法是采用“定时批量计算 前端展示最新结果”的模式。7.1 手动触发最简单的模式是把 Spark 计算过程写成一个独立脚本run_recommend.py在需要更新推荐结果时手工执行python run_recommend.py执行完成后再执行一下 Django 的数据同步脚本把结果导入 MySQL。这种模式适合答辩前演示也适合新手理解整个数据流。7.2 自动定时触发在 Linux 环境下可以使用crontab做每日定时任务0 2 * * * cd /path/to/project python run_recommend.py logs/recommend.log 21在 Windows 环境下可以使用“任务计划程序”指向 Python 脚本。定时任务的优点是稳定缺点是如果中间某个环节报错你可能要到第二天才发现。所以脚本里一定要保留完整日志输出。7.3 并发和生产化思路如果是自己的学习项目并发问题不太明显。但如果部署到服务器上就需要考虑Spark 任务和 Django 服务不要在同一端口或同一进程里跑。Spark 任务的内存限制要设置好避免把服务器内存占满。数据库连接池要配置合理避免高并发查询时连接数耗尽。推荐结果表可以按日期分区保留最近几天的结果避免表无限膨胀。这些点不用在毕设答辩时全部展开但可以作为“后续优化方向”提前写进论文里。8. 常见报错与排查链路按顺序查能省掉很多无效操作这个项目涉及的组件多报错样式也五花八门。根据过往经验我把问题排查顺序固定成下面这条链路先看现象是启动失败、日志报错、页面无数据、还是推荐结果一直不变再看输入CSV 文件路径对不对、字段名是否完全一致、编码是否为 UTF-8、有没有空行或缺失值。再看环境Hadoop 是否正常启动、Spark 能否提交任务、端口是否被占用、Java 版本是否匹配、MySQL 服务是否在运行。再看参数Spark 任务内存是否设置过低、Django 数据库连接参数是否正确、并发数是否过大。最后看代码逻辑是不是 join 字段写错了是不是过滤条件把结果全部过滤掉了。下面说几个高频问题。8.1 Spark 读取 HDFS 文件报错常见原因是文件路径写错或者文件根本没上传成功。先用hdfs dfs -ls /recommend/input/确认文件存在再检查 Spark 读取时使用的主机名和端口是否与 Hadoop 配置一致。如果在本地模式用 9000 端口而 Hadoop 配置的是 8020就会报连接失败。8.2 Django 页面能打开但推荐列表为空这是最常遇到的情况。原因通常不是推荐逻辑错了而是推荐结果没有写入数据库或者写入的是另一个用户的数据。排查顺序先确认当前登录用户是谁。再在 MySQL 里查一下RecommendResult表有没有该用户的记录。如果没有就重新执行一次 Spark 结果导入脚本并查看日志。如果记录存在但还是空检查模板里的字段名是否正确。8.3 Spark 任务执行很慢如果你用伪分布式或本地模式处理几十万条以上数据慢是正常的。可以先用.limit(1000)或采样确认逻辑正确再全量跑。另外关闭不必要的日志输出减少 shuffle 分区数也会有一定效果。8.4 Hadoop 无法正常启动 DataNode常见原因是格式化后多次重新启动导致 NameNode 和 DataNode 的集群 ID 不一致。解决方法是先停止 Hadoop 相关进程删除临时目录下的数据文件再重新格式化。注意删除数据目录前要确认是测试环境否则会丢失数据。9. 答辩和论文里可以重点突出的“工程亮点”毕设项目做出来了能跑通只是基础分。要想拿高分还需要在论文和答辩 PPT 中突出几个工程化能力9.1 数据分层设计你可以把项目的数据划分为原始行为层、特征统计层、推荐结果层。每一层对应不同数据表或 HDFS 路径。答辩时可以画一张三层架构图解释为什么需要这么分层原始层保留日志完整性特征层方便后续复用结果层面向业务展示。这一句话就能体现出你对数据工程的理解而不是只会调框架。9.2 推荐解释机制在推荐结果表里加入reason字段让前端可以展示“因为你对科技类视频感兴趣所以为你推荐了这条”。这个细节能拉开差距因为很多同学的推荐系统就是一个黑盒看不到推荐原因。能解释为什么推荐说明你考虑了推荐系统的可信度问题。9.3 行为模拟与实验验证你可以设计一个小实验在模拟数据里让某个用户大量点击美食类视频然后重新跑推荐任务观察首页推荐流是否向美食类倾斜。如果确实倾斜了就说明推荐逻辑生效了。这种“预设条件 → 执行计算 → 验证结果”的实验过程是论文里很好写的部分也能证明项目不是摆设。9.4 异常数据处理比如用户观看时长超过视频总时长、同一用户同一视频在 1 秒内重复点击、空值和脏数据。你在 Spark 清洗阶段如何处理这些问题都可以单独写一小节。这些内容最能体现“真实工程能力”哪怕只是做了很基础的去重和阈值过滤也值得写。10. 经验总结哪些地方不用过度追求哪些地方必须一步到位最后说几条我用实际踩坑换来的经验。第一环境版本尽量一开始就固定。Hadoop、Spark、JDK、Python 版本之间互相影响。不要用最新版冒险也不要混装多个大版本。建议用一套稳定组合后把版本号写进项目 README。第二数据量不要一开始就堆太大。我见过很多同学准备生成几十万条数据结果 Spark 任务跑得很慢每次调试都要等好久。建议先用 2 万条数据把流程跑通确认结果正确后再考虑是否扩大。第三不要尝试让 Django 主页实时触发 Spark 计算。这个方案在并发高的时候一定会卡死。把“计算”和“展示”解耦开才是这个项目应该体现出的架构思维。第四日志一定要早写。每跑一步都写日志或者至少保留命令行输出。排查问题的时候日志是最可靠的依据比肉眼检查代码高效得多。第五答辩演示前一定把流程顺序固定下来。先启动 Hadoop 相关服务再启动 MySQL再启动 Django最后跑一次 Spark 计算并展示最新推荐结果。每一步的预期结果都要提前确认一次不要在答辩现场临时改参数或重新生成数据。这个项目最大价值不是“做出来一个网站”而是让你完整经历一遍从数据产生、存储、计算到展示的全链路开发。只要这条链路能讲清楚代码哪怕简单一些也是一个合格的毕设项目。
返回列表