基于Spark的个性化短视频推荐系统设计与实现
这次我们来看一个计算机毕业设计方向的项目基于 Spark 的个性化短视频推荐系统。技术栈是 Python、Spark、Hadoop、Django典型的“大数据存储 分布式计算 Web 应用”组合。这类项目在毕业设计里很常见但完整把推荐算法、离线计算、在线服务和后台管理串起来的并不多。这篇文章会把系统架构、核心功能、推荐流程、部署步骤、接口设计和常见坑都拆开讲清楚方便你直接参考或二次开发。先给结论这个项目适合正在准备 Python 方向毕业设计、想接触大数据技术栈、又不想完全脱离 Web 开发的同学。核心价值在于它打通了从数据采集、用户行为日志、离线推荐计算到前端短视频展示的完整链路。对比单纯的 Django 管理系统或单纯的 Spark 算法Demo这个题目的技术覆盖面更宽也更容易在答辩时展开讲。文章会按下面几个部分展开核心能力速览、系统架构设计、功能模块拆解、环境准备与部署、推荐算法实现思路、接口与批量任务、性能观察、常见问题排查、最佳实践与合规提醒。1. 核心能力速览能力项说明项目类型毕业设计 / 课程设计完整项目后端框架Django数据存储MySQL HDFSHadoop 分布式文件系统离线计算Spark 批处理负责用户行为日志清洗和推荐结果计算推荐策略基于物品的协同过滤 / 基于用户的协同过滤 / 热度补冷启动可按数据情况调整Web 管理端Django Admin 自定义后台管理用户、短视频、分类、推荐结果前端展示短视频信息流页面支持分类筛选、热门列表、个性化推荐列表批量任务Spark 离线任务可定时运行支持全量用户推荐结果更新API 支持Django REST Framework 风格接口可扩展在线推荐查询部署方式本机开发环境 / Linux 服务器部署Hadoop 伪分布式或集群适合人群计算机相关专业毕设、大数据课程设计、想入门推荐系统的开发者从材料看项目附带源码、文档报告和代码讲解适合作为毕设主体项目。启动时可以拆成“先跑推荐计算再启动 Web 服务”两条线。推荐结果可以先离线算好写入 MySQL前端直接查询展示避免在线推荐接口在低配置机器上响应过慢。2. 适用场景与使用边界这类项目最适合两类人。第一类是准备做毕业设计的本科生。题目本身把大数据和 Web 开发结合起来既有算法模型可以实现又有可视化页面可以展示答辩时“系统功能 算法原理 实验对比”都能撑起来。相比纯管理系统推荐系统更容易讲出技术含量。第二类是入门大数据开发的开发者。很多人学完 Spark 和 Hadoop 以后不知道能做什么完整项目这个系统提供了一个从数据产生到数据计算再到数据消费的闭环案例。通过修改 Spark 任务里的推荐逻辑就能验证不同推荐算法在真实行为数据上的效果。使用边界也要说清楚。这个项目的定位是学习和教学演示不是工业级推荐平台。它不会覆盖实时特征计算、向量召回、深度排序模型、A/B Test 等生产环境能力。行为数据量级在百万以下时单机伪分布式 Hadoop Spark 完全可以支撑如果数据量达到亿级就需要考虑集群部署和更复杂的推荐架构。另外涉及用户行为数据采集和短视频内容展示时必须注意数据合规。如果数据来自公开数据集要确认数据集的使用许可如果自己爬取视频数据要遵守平台服务条款不能抓取用户个人信息。项目用于毕设演示时尽量使用脱敏数据或自行构造的测试数据。3. 系统架构设计整个系统可以拆成数据层、计算层、应用层三个部分。数据层负责存储三类数据用户基本信息、短视频元数据、用户行为日志。短视频元数据和用户信息放在 MySQL方便 Django 直接查询用户行为日志可以写入 HDFS供 Spark 离线任务读取。这样的好处是 Web 业务与大数计算之间的存储隔离互不影响。计算层是 Spark 离线任务。任务逻辑大致如下从 HDFS 读取用户行为日志 - 数据清洗 - 生成用户-物品评分矩阵 - 调用推荐算法计算相似度或推荐结果 - 把结果写入 MySQL 推荐结果表。Spark 在这里承担了“离线计算引擎”的角色。如果暂时没有 Hadoop 环境也可以先用本地文件模拟 HDFS 数据源先跑通推荐流程再切换数据源。应用层是 Django 项目负责面向用户的短视频展示、分类浏览、热门排行、个性化推荐等页面同时提供管理员后台。Django 模型操作 MySQL通过 ORM 读取推荐结果表中的数据再渲染到前端页面。架构上比较关键的设计是“推荐结果落库”。在线推荐服务如果每次请求都实时计算性能压力很大。更稳妥的做法是 Spark 定时计算一次全量推荐结构写入 MySQLDjango 页面直接查询。这样 Web 端响应速度快Spark 和 Django 两个模块的耦合度也低。4. 功能模块拆解4.1 用户模块用户模块承担注册、登录、个人信息管理功能。Django 自带的用户体系可以直接使用也可以继承 AbstractUser 扩展字段比如用户偏好标签、注册时间、最近登录时间等。行为数据采集依赖用户标识所以用户登录后需要把 user_id 写入会话或前端存储后续操作短视频时带上该标识。4.2 短视频内容管理模块短视频元数据包括视频标题、视频封面、视频链接、分类、上传时间、播放量、点赞数、评论数等。这部分在 Django 里就是常规的模型定义和管理后台。为了演示效果可以直接使用在线视频链接或本地静图替换本质是操作一条视频记录前端展示对应字段。4.3 用户行为采集模块推荐系统依赖用户行为数据。每一次曝光、点击、播放、点赞、收藏、评论、分享都可以记录为一条行为日志。工程实现上有两种路径直接写入 MySQL 行为表写入本地日志文件定时同步到 HDFS第二种更接近真实场景。Django 接收用户操作后先落日志文件再由日志采集脚本上传到 HDFSSpark 定时读取。如果实验环境有限可以直接写 MySQLSpark 从 MySQL 读取行为数据。这个取舍不会影响推荐算法的验证。4.4 推荐计算模块推荐算法是系统核心。从实现难度和毕设答辩效果看推荐算法可以选择以下组合热度推荐按播放量、点赞量、时间衰减计算热度分适合冷启动。基于物品的协同过滤计算物品相似度根据用户历史行为推荐相似物品。基于用户的协同过滤找到相似用户推荐相似用户喜欢的物品。混合推荐将多种算法结果加权融合提高推荐覆盖率和多样性。Spark 在推荐计算中有天然优势RDD 和 DataFrame 适合处理大规模用户行为数据mlib 库提供了 ALS 矩阵分解算法可以用在协同过滤场景。对于毕设项目基于物品的协同过滤实现起来更直观也更容易讲清楚相似度计算公式和推荐流程。4.5 前端信息流模块前端展示短视频信息流按推荐结果排序。实现方式可以很简单Django 视图函数从推荐结果表读取当前用户的推荐视频列表传给模板渲染。页面支持上滑加载更多对应分页查询。如果希望效果更接近短视频 App可以引入 Vue 或 jQuery 做异步加载但这会额外增加前后端联调成本。建议第一次开发时先用 Django 模板 分页实现保证核心链路稳定后再优化交互。5. 环境准备与前置条件这是一个 Django Spark Hadoop 组合项目环境准备要分两层来看。Web 开发层需要 Python、Django、MySQL 和必要的 Python 依赖库。推荐使用 Python 3.8 或 3.9 版本Django 使用 3.x 或 4.x 的常青版本。太新的 Python 版本可能遇到部分依赖库编译问题太老的版本又不利于后续扩展。大数据层需要 JDK、Hadoop、Spark。这里给一个通用检查清单组件作用版本建议JDKHadoop 和 Spark 运行基础JDK 8 或 JDK 11HadoopHDFS 存储与资源调度Hadoop 3.xSpark离线推荐计算Spark 3.x需匹配 Scala 版本MySQL业务数据存储MySQL 5.7 或 8.0PythonDjango 开发Python 3.8DjangoWeb 后端Django 3.2没有现成大数据环境的同学可以先用本地模式跑 Spark 任务。安装 Spark 后不启动 Hadoop 集群直接把 Spark 作为本地计算引擎使用核心目的是先跑通推荐链路的计算逻辑。后续再根据实验需要切换到 HDFS 数据源。硬盘空间至少预留 20GBHadoop 和 Spark 解压后大约占用 3 到 5GB再加上 Python 依赖和系统缓存空间宽松一点更稳妥。内存方面8GB 内存跑单机伪分布式可以工作16GB 更流畅。6. 安装部署与启动方式6.1 Django 后端启动先创建 Python 虚拟环境并安装依赖python -m venv venv source venv/bin/activate # Windows 使用 venv\Scripts\activate pip install django3.2.* mysqlclient配置数据库连接在 Django 的 settings.py 中修改DATABASES { default: { ENGINE: django.db.backends.mysql, NAME: video_recommend, USER: root, PASSWORD: your_password, HOST: 127.0.0.1, PORT: 3306, } }执行数据库迁移并启动服务python manage.py makemigrations python manage.py migrate python manage.py createsuperuser python manage.py runserver 0.0.0.0:8000浏览器访问http://127.0.0.1:8000能看到系统首页和用户登录入口说明 Django 部分启动成功。6.2 Hadoop 伪分布式搭建Hadoop 伪分布式是单机模拟分布式环境。需要配置 core-site.xml、hdfs-site.xml、yarn-site.xml 三个核心文件然后格式化 NameNodehdfs namenode -format start-dfs.sh start-yarn.sh jps启动后看到 NameNode、DataNode、ResourceManager、NodeManager 进程就说明 Hadoop 环境正常。部分电脑可能因为主机名映射或 SSH 配置问题启动失败需要确认/etc/hosts中配置了主机名对应 127.0.0.1。6.3 Spark 推荐任务启动Spark 任务可以打包成 Python 脚本使用 spark-submit 提交。下面是一个简化的任务入口示例spark-submit \ --master local[2] \ --driver-memory 2g \ recommend_job.py \ --input hdfs://localhost:9000/user/logs \ --output jdbc:mysql://127.0.0.1:3306/video_recommend如果暂时没有 HDFS可以用本地文件作为输入spark-submit \ --master local[2] \ recommend_job.py \ --input ./data/user_behavior.log \ --output jdbc:mysql://127.0.0.1:3306/video_recommendrecommend_job.py 内部需要做四件事读取行为日志、构造用户物品矩阵、计算推荐结果、写入 MySQL。这里给一个基于物品协同过滤的伪代码示例from pyspark.sql import SparkSession from pyspark.sql.functions import col spark SparkSession.builder \ .appName(VideoRecommendJob) \ .getOrCreate() # 1. 读取行为数据 df spark.read.csv(./data/user_behavior.log, headerTrue) # 字段示例user_id, video_id, behavior_type, timestamp # 2. 构造评分比如播放1点赞3收藏5 df df.withColumn(score, col(behavior_score)) # 3. 按用户和视频聚合评分矩阵 rating_matrix df.groupBy(user_id, video_id).avg(score) # 4. 计算视频相似度 # 核心是余弦相似度或 Jaccard 相似度具体实现可按数据量选择 # 这里省略中间计算细节最终生成 video_sim 表或 DataFrame # 5. 为每个用户生成 Top N 推荐并写入 MySQL result_df generate_recommendations(rating_matrix, video_sim) result_df.write \ .format(jdbc) \ .option(url, jdbc:mysql://127.0.0.1:3306/video_recommend) \ .option(dbtable, recommend_result) \ .option(user, root) \ .option(password, your_password) \ .mode(overwrite) \ .save()需要注意这一段是通用模板真实的相似度计算函数需要根据你的评分矩阵和推荐算法目标去实现。毕业设计中基于物品协同过滤的相似度计算通常使用余弦相似度计算公式为两个物品向量夹角的余弦值。Spark 里可以把视频向量 RDD 化再用 cartesian 或其他算子计算但要注意数据量避免 OOM。7. 推荐算法实现思路7.1 冷启动策略新用户没有行为数据协同过滤算不出结果。最常用的方案是用热度推荐兜底。热度分可以综合播放量、点赞量、收藏量并加入时间衰减因子避免旧视频长期霸榜。伪代码逻辑如下hot_score play_count * 1 like_count * 3 collect_count * 5 # 按天衰减 from datetime import datetime, timedelta days_gap (datetime.now() - publish_time).days final_hot_score hot_score / (1 days_gap * 0.1)新视频没有行为数据除了时间因子还可以考虑分类匹配。用户注册时选择兴趣分类新视频按分类权重优先展示。7.2 基于物品的协同过滤核心思路是“你喜欢视频 A那么和 A 相似的视频 B 也值得推荐”。第一步构建视频评分向量。每个视频被用户观看、点赞、收藏后形成一个user_id - score的向量。第二步计算视频两两之间的相似度。常用的余弦相似度公式是similarity(A, B) sum(A_i * B_i) / (sqrt(sum(A_i^2)) * sqrt(sum(B_i^2)))第三步对用户已产生过行为的视频集合取其相似视频并按相似度与用户评分的加权值排序得到推荐候选集。这种方法在用户行为稀疏时效果会受限所以需要结合热度推荐和分类推荐做混合。7.3 基于 ALS 的矩阵分解Spark mllib 提供了 ALS 算法可以直接用在协同过滤场景。训练输入是(user_id, video_id, rating)三元组输出是用户因子矩阵和物品因子矩阵。预测评分时计算用户向量和物品向量的点积。from pyspark.ml.recommendation import ALS als ALS( userColuser_id, itemColvideo_id, ratingColscore, coldStartStrategydrop ) model als.fit(training_data) predictions model.transform(test_data)ALS 的调参重点是 rank因子数量、iterations迭代次数、regParam正则化系数。毕设阶段可以先固定一组参数跑通推荐流程后再做小规模网格搜索对比。8. 接口 API 与批量任务8.1 推荐结果查询接口Django 可以暴露一个简单 JSON 接口方便前端异步加载推荐列表。示例视图如下from django.http import JsonResponse from .models import RecommendResult, Video def recommend_list(request, user_id): recommend_ids ( RecommendResult.objects .filter(user_iduser_id) .order_by(-score) .values_list(video_id, flatTrue)[:30] ) videos Video.objects.filter(id__inlist(recommend_ids)) data [ { video_id: v.id, title: v.title, cover: v.cover_url, video_url: v.video_url, play_count: v.play_count, } for v in videos ] return JsonResponse({code: 0, data: data})前端通过fetch(/api/recommend/1/)获取数据再渲染到页面。这样可以避免模板硬编码也让前后端职责更清晰。8.2 行为上报接口用户点击、播放、点赞等行为可以通过 POST 接口上报import json import logging from django.views.decorators.csrf import csrf_exempt from django.http import JsonResponse logger logging.getLogger(behavior) csrf_exempt def behavior_report(request): if request.method POST: body json.loads(request.body) log_line {user_id}|{video_id}|{behavior}|{timestamp}.format( user_idbody.get(user_id), video_idbody.get(video_id), behaviorbody.get(behavior), timestampbody.get(timestamp) ) logger.info(log_line) return JsonResponse({code: 0, msg: ok}) return JsonResponse({code: -1, msg: method not allowed})行为日志落到本地文件后可以通过脚本定时上传到 HDFS。如果只想在 MySQL 里完成闭环也可以直接 ORM 写入行为表。8.3 Spark 批量任务调度最简单的批量调度方式是用 crontab 定时执行 spark-submit# 每天凌晨 2 点执行推荐计算 0 2 * * * /opt/spark/bin/spark-submit \ --master local[2] \ /opt/video_recommend/recommend_job.py \ --input hdfs://localhost:9000/user/logs \ --output jdbc:mysql://127.0.0.1:3306/video_recommend批量任务的关键是要有幂等性。推荐结果表每次全量覆盖写入使用mode(overwrite)可以避免重复数据累积。如果推荐结果分批次写入建议增加批次号字段前端只读取最新批次。涉及批量处理时要考虑任务失败重试。最简单的方式是在 shell 脚本里判断退出码失败时发送告警或延长等待时间后重跑if [ $? -ne 0 ]; then echo recommend job failed at $(date) /var/log/recommend_job.log exit 1 fi9. 资源占用与性能观察9.1 查看 Spark 任务资源占用Spark Web UI 默认在http://localhost:4040可以观察每个 Stage 耗时、输入数据量、Shuffle 大小和 Executor 内存。推荐任务跑完后重点看两个指标Shuffle Read/Write 大小如果很大说明 join 操作涉及大量数据迁移。Executor 内存使用如果 GC 频繁或 OOM需要调大 driver-memory 或减少分区数。9.2 MySQL 和 Django 侧观察Django 开发服务器下页面响应时间会因为推荐结果表的数据量而不同。如果推荐结果表有索引例如(user_id, score)联合索引查询会快很多。ALTER TABLE recommend_result ADD INDEX idx_user_score (user_id, score);前端页面加载慢时优先检查是否执行了 N1 查询。推荐列表取回来后关联查询 Video 表要使用select_related或一次性filter(id__in...)避免在循环里反复查数据库。9.3 降低资源消耗低配机器下可以从几个方向压降资源Spark 任务采用local[2]而不是local[*]限制并发线程数。行为数据先按月或按天分区只对最近一个月数据做全量计算。推荐结果只保留 Top 50不把全量候选入库。Django 使用缓存接口缓存推荐列表结果例如 cache_page 或 Redis 缓存。10. 常见问题与排查方法问题现象可能原因排查方式解决方案Django 启动报数据库连接错误MySQL 未启动或密码错误检查 MySQL 服务状态用命令行连接数据库启动 MySQL修改 settings.py 数据库配置create_superuser 无法创建数据库迁移未执行查看 migrate 是否执行成功执行 makemigrations 后再执行 migrateSpark 任务报 ClassNotFoundSpark 与 Scala/JDK 版本不匹配查看 Spark 版本说明重新下载匹配 Spark 发行版确认 JDK 版本Spark 读取 HDFS 路径失败Hadoop 未启动或路径不对hdfs dfs -ls /user/logs 检查目录启动 Hadoop 或先使用本地文件测试spark-submit 提交后很快失败Python 依赖在 Spark executor 中缺失查看 yarn 日志或 driver 日志用 --py-files 上传依赖包或改用 pandas 兼容逻辑推荐结果为空行为数据太少或冷启动未处理检查 rating_matrix 是否有数据补充测试行为数据冷启动阶段返回热度推荐Web 页面响应慢推荐结果表无索引EXPLAIN 查看查询计划添加联合索引减少关联查询次数行为日志一直增长日志轮转缺失查看磁盘占用配置 logrotate 或定时清理脚本推荐结果重复多次运行批任务且未覆盖写入检查 recommend_result 表数据量使用 overwrite 模式或增加批次号11. 最佳实践与使用建议第一次跑通项目时不要急着调算法指标。先按“MySQL 导入测试数据 - Django 启动后台 - 构造行为日志 - 运行 Spark 推荐任务 - 查看推荐结果”的顺序走一遍主链路。主链路通了再逐步换真实数据集或调整推荐策略。工程上建议做下面几件事保留一套最小可运行配置。比如使用本地文件模拟 HDFS、使用 SQLite 临时替代 MySQL确保在没装 Hadoop 的机器上也能源码演示。数据目录分清楚。原始行为数据、清洗后数据、推荐输出结果、日志文件分别放在独立目录。推荐任务增加日志输出。每一步处理的数据量、耗时、异常记录都写入日志。批量任务要控制幂等性。每次覆盖写结果表或使用批次号避免重复推荐。接口服务如果监听在公网端口要限制 IP 访问范围避免被滥用。涉及用户信息、行为数据时要脱敏处理。展示推荐结果时不能泄露他人隐私。关于版权和合规这里要多说几句。如果系统里使用了真实短视频平台的数据或素材必须确认来源合法不能把抓取的视频直接用于公开演示。推荐系统处理用户行为日志时要明确告知用户数据用途不能收集非必要的个人隐私信息。毕业设计公开答辩或上传论文前最好使用自行构造的测试数据。12. 总结与下一步这个项目最值得尝试的是完整跑通“Hadoop 存储日志数据 - Spark 离线计算推荐结果 - Django 展示短视频信息流”的闭环。相比单一技术栈的毕设它让你在同一个项目里接触到 Web 开发、大数据存储、分布式计算和推荐算法四个方向后续写简历、准备面试也都有素材可以讲。最先验证的功能应该是 Spark 推荐任务能否把结果成功写入 MySQL。这一条通了整个项目的核心价值就立住了。最容易踩的坑有两处一是环境和依赖版本不匹配Hadoop、Spark、JDK、Python 版本稍有不一致就容易启动失败二是推荐任务跑完但 MySQL 里的结果表没有数据通常是因为 JDBC 写入依赖包缺失需要把 MySQL Connector/J 放到 Spark 的 jars 目录。后续可以继续扩展的方向包括引入 Redis 缓存在线推荐结果、增加基于 ALS 的矩阵分解对比实验、加入用户实时行为流处理如 Spark Streaming、把 Django 前端替换成 Vue 单页应用。无论走哪个方向现有项目已经提供了一个可稳定扩展的地基值得在里面继续深挖。