基于 Spark 的个性化短视频推荐系统,算是我见过技术栈比较完整的 Python 毕设项目之一。它没有把推荐做成一个“demo 级”的摆设,而是把用户行为采集、Hadoop 日志存储、Spark 离线计算、Django Web 服务完整串成了一条流水线。如果你正在选毕设题目,或者想自己动手搭一套推荐系统做练习,这篇拆解可以直接作为参考路线。
这个项目最大的特点是:用了 Spark + Hadoop 做大数据处理,用 Django 做业务后端,前端负责视频浏览、播放、点赞、收藏这些交互。也就是说,你可以在一个项目里同时体现“大数据计算”和“Web 开发”两种能力,答辩时技术点很密,项目含金量也会比单纯的 CRUD 系统高不少。本文会从系统架构、推荐算法设计、环境准备、部署启动、功能验证、性能观察、常见问题几个维度完整拆解,末尾还会给出代码示例和毕设避坑建议。
1. 核心能力速览
| 能力项 | 说明 |
|---|---|
| 项目类型 | Python 毕业设计 / 大数据推荐系统 |
| 技术栈 | Python、Spark、Hadoop、Django、MySQL |
| 核心功能 | 用户注册登录、视频浏览播放、点赞收藏评论、个性化推荐、后台管理 |
| 推荐算法 | 基于用户的协同过滤、基于物品的协同过滤、基于内容的推荐、热度兜底 |
| 数据存储 | MySQL 存业务数据,HDFS 存日志与中间结果 |
| 计算框架 | Spark 离线批量计算 |
| Web 框架 | Django + 前端模板 |
| 运行环境 | Linux 推荐,Windows 也可做开发调试 |
| 扩展能力 | 可接 Redis 缓存、Celery 异步任务、定时调度 |
| 适用场景 | 毕业设计、课程设计、推荐系统入门练习 |
从这些能力项能看出来,这个项目的核心价值不是某一个算法多前沿,而是完整覆盖了“数据采集 → 数据存储 → 离线计算 → 推荐结果 → Web 展示”的闭环。对毕设来说,这个闭环比单纯调一个推荐算法库更有说服力。
2. 系统架构与推荐流程
2.1 分层架构
整个系统可以拆成四层:
- 数据采集层:前端页面埋点,用户浏览、播放、点赞、收藏、评论视频时,行为数据通过接口上报,生成日志。
- 数据存储层:业务数据(用户表、视频表、行为表、推荐结果表)存 MySQL;原始日志和 Spark 计算的中间结果存 HDFS。
- 计算层:Spark 定时读取用户行为日志,执行数据清洗、特征统计、相似度计算、推荐列表生成。
- 应用层:Django 提供 Web 服务,处理用户请求,从数据库读取推荐结果并渲染页面。
理论上这就是一个典型的离线推荐系统架构。实时推荐可以做,但那是加分项,毕设核心先把离线链路跑通。
2.2 推荐数据流向
一次完整的推荐流程如下:
- 用户访问网站,浏览或播放视频。
- 前端把行为数据发送到 Django 接口。
- Django 将行为写入 MySQL 业务库,同时把日志写入 HDFS。
- Spark 离线任务定时执行,读取 HDFS 日志,做数据清洗和推荐计算。
- Spark 生成每个用户的 Top-N 推荐列表,写回 MySQL 推荐结果表。
- 用户刷新首页,Django 从推荐结果表取出数据,渲染个性化视频列表。
这套流程的关键点是:推荐结果不是实时算的,而是离线算好后写库,Web 端直接查库展示。这样设计对毕设项目非常合理,既避免实时计算的高复杂度,又能把推荐效果清晰展示出来。
3. 推荐算法设计与实现
3.1 算法选型
推荐算法部分,项目主要围绕以下几类展开:
- 基于用户的协同过滤(UserCF):找到与当前用户兴趣相似的其他用户,把这些用户喜欢的视频推荐给当前用户。
- 基于物品的协同过滤(ItemCF):找到与用户历史上喜欢过的视频相似的视频,推荐给当前用户。
- 基于内容的推荐(Content-based):根据视频的分类、标签、标题关键词等属性,推荐同类型视频。
- 热度推荐与新视频推荐:解决冷启动问题,对没有行为记录的新用户,推荐全局热门视频或最新上传视频。
其中 UserCF 和 ItemCF 是推荐系统的两大基础算法,也是答辩时最容易被问到的点。建议把协同过滤的原理、相似度公式、为什么会有冷启动问题都弄清楚。
3.2 用户行为加权
不同行为对用户兴趣的贡献不一样。Spark 离线统计时,可以给行为赋予不同权重:
| 行为类型 | 权重 |
|---|---|
| 播放 | 1 |
| 点赞 | 2 |
| 收藏 | 3 |
| 评论 | 2 |
| 分享 | 3 |
代码示例(Spark SQL 统计行为得分):
SELECT user_id, video_id, SUM( CASE behavior_type WHEN 'play' THEN 1 WHEN 'like' THEN 2 WHEN 'favorite' THEN 3 WHEN 'comment' THEN 2 WHEN 'share' THEN 3 ELSE 0 END ) AS score FROM user_behavior_log GROUP BY user_id, video_id3.3 协同过滤相似度计算
协同过滤的核心是相似度计算。以 UserCF 为例:
- 构造用户-物品评分矩阵。
- 计算用户之间的相似度,常用余弦相似度或皮尔逊相关系数。
- 选取 Top-K 相似用户。
- 汇总相似用户喜欢的视频,排除当前用户已看过的,计算推荐得分。
Python 伪代码示例如下:
import math from collections import defaultdict def cosine_similarity(user_vector1, user_vector2): """计算两个用户向量之间的余弦相似度""" common_items = set(user_vector1.keys()) & set(user_vector2.keys()) if not common_items: return 0.0 dot_product = sum(user_vector1[item] * user_vector2[item] for item in common_items) norm1 = math.sqrt(sum(value ** 2 for value in user_vector1.values())) norm2 = math.sqrt(sum(value ** 2 for value in user_vector2.values())) if norm1 == 0 or norm2 == 0: return 0.0 return dot_product / (norm1 * norm2) def user_based_recommend(user_id, user_item_matrix, top_k=10): """基于用户的协同过滤推荐""" target_user_vector = user_item_matrix[user_id] similarity_scores = [] for other_user_id, other_vector in user_item_matrix.items(): if other_user_id == user_id: continue sim = cosine_similarity(target_user_vector, other_vector) similarity_scores.append((other_user_id, sim)) # 排序取前 K 个相似用户 similarity_scores.sort(key=lambda x: x[1], reverse=True) top_k_users = similarity_scores[:top_k] # 候选物品得分 candidate_score = defaultdict(float) for other_user_id, sim in top_k_users: for video_id, score in user_item_matrix[other_user_id].items(): if video_id in target_user_vector: continue candidate_score[video_id] += sim * score # 排序返回推荐列表 sorted_candidates = sorted(candidate_score.items(), key=lambda x: x[1], reverse=True) return [video_id for video_id, _ in sorted_candidates]这个逻辑在 Spark 里实现时,可以把用户-物品矩阵转换成 RDD 或 DataFrame,用 join 和 groupBy 来替代循环。数据量大时,基于 Spark 的分布式计算优势就体现出来了。
4. 数据采集与数据库设计
4.1 数据库表设计
项目涉及的 MySQL 核心表大致包括:
- 用户表(user):用户 ID、用户名、密码、头像、注册时间。
- 视频表(video):视频 ID、标题、封面、视频地址、分类、标签、上传者、上传时间。
- 用户行为表(user_behavior):行为 ID、用户 ID、视频 ID、行为类型、创建时间。
- 推荐结果表(recommend_result):推荐 ID、用户 ID、视频 ID、推荐得分、生成时间。
Django 模型定义示例:
from django.db import models from django.contrib.auth.models import AbstractUser class User(AbstractUser): avatar = models.URLField(blank=True, null=True, verbose_name="头像") created_at = models.DateTimeField(auto_now_add=True, verbose_name="注册时间") class Meta: db_table = "user" verbose_name = "用户" class Video(models.Model): title = models.CharField(max_length=200, verbose_name="标题") cover_url = models.URLField(verbose_name="封面地址") video_url = models.URLField(verbose_name="视频地址") category = models.CharField(max_length=50, verbose_name="分类") tags = models.CharField(max_length=200, verbose_name="标签") uploader = models.ForeignKey(User, on_delete=models.CASCADE, verbose_name="上传者") created_at = models.DateTimeField(auto_now_add=True, verbose_name="上传时间") class Meta: db_table = "video" verbose_name = "视频" class UserBehavior(models.Model): BEHAVIOR_CHOICES = ( ("play", "播放"), ("like", "点赞"), ("favorite", "收藏"), ("comment", "评论"), ("share", "分享"), ) user = models.ForeignKey(User, on_delete=models.CASCADE, verbose_name="用户") video = models.ForeignKey(Video, on_delete=models.CASCADE, verbose_name="视频") behavior_type = models.CharField(max_length=20, choices=BEHAVIOR_CHOICES, verbose_name="行为类型") created_at = models.DateTimeField(auto_now_add=True, verbose_name="行为时间") class Meta: db_table = "user_behavior" verbose_name = "用户行为"4.2 行为日志上报
用户产生行为时,前端通过 Ajax 或埋点脚本上报到 Django 接口。Django 视图收到请求后,一方面写入 MySQL,另一方面把日志格式化为一行 JSON,追加写入 HDFS 或本地日志目录。
import json import logging from django.http import JsonResponse from django.views.decorators.csrf import csrf_exempt from .models import UserBehavior logger = logging.getLogger("recommend") @csrf_exempt def report_behavior(request): if request.method == "POST": data = json.loads(request.body) user_id = data.get("user_id") video_id = data.get("video_id") behavior_type = data.get("behavior_type") # 写入 MySQL UserBehavior.objects.create( user_id=user_id, video_id=video_id, behavior_type=behavior_type ) # 写日志,供 Spark 离线消费 log_line = json.dumps({ "user_id": user_id, "video_id": video_id, "behavior_type": behavior_type, "timestamp": datetime.now().isoformat() }) logger.info(log_line) return JsonResponse({"code": 0, "message": "success"})这里注意一个细节:MySQL 表里的行为数据可以做实时展示,比如“我点赞过的视频”;HDFS/日志文件里的数据才是 Spark 推荐计算的输入。两者分开,职责更清晰。
5. 环境准备与前置条件
这个项目的环境搭建主要涉及 JDK、Hadoop、Spark、MySQL、Python 和 Django。下面给出一份通用检查清单,具体版本号要根据你自己下载的安装包确定,不要盲目照抄网上命令。
5.1 软件清单
| 软件 | 用途 | 说明 |
|---|---|---|
| JDK 1.8 | Hadoop 和 Spark 运行依赖 | 必须安装,并配置 JAVA_HOME |
| Hadoop | 分布式存储(HDFS) | 开发环境可单机伪分布式部署 |
| Spark | 离线计算引擎 | 依赖 Hadoop,负责跑推荐任务 |
| MySQL | 业务数据库 | 存用户、视频、行为、推荐结果 |
| Python 3 | Django Web 开发 | 建议用虚拟环境管理依赖 |
| Django | Web 框架 | 安装 django、pymysql 等依赖 |
| d3.js / ECharts | 数据可视化(可选) | 如果系统里要做统计图表 |
5.2 环境变量配置
Linux 环境下,需要在~/.bashrc或/etc/profile中配置环境变量:
export JAVA_HOME=/usr/local/jdk1.8 export HADOOP_HOME=/usr/local/hadoop export SPARK_HOME=/usr/local/spark export PATH=$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$SPARK_HOME/bin配置完成后执行:
source ~/.bashrc java -version hadoop version spark-shell --version三个命令都能正确输出版本号,说明基础环境没问题。
5.3 Python 虚拟环境
推荐使用 venv 或 conda 创建独立环境,避免污染系统 Python:
python3 -m venv venv source venv/bin/activate pip install django pymysql requests6. 安装部署与启动方式
6.1 Hadoop 伪分布式部署
开发学习阶段不需要搭建真实集群,单机伪分布式模式足够跑通流程。
先修改$HADOOP_HOME/etc/hadoop/core-site.xml:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/usr/local/hadoop/tmp</value> </property> </configuration>再修改hdfs-site.xml,设置副本数为 1:
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> </configuration>然后执行启动命令:
cd $HADOOP_HOME bin/hdfs namenode -format sbin/start-dfs.sh启动后用jps查看进程,能看到 NameNode、DataNode、SecondaryNameNode 三个进程,说明 HDFS 启动成功。
6.2 Spark 离线任务提交
推荐任务写好后,打成 Python 脚本或 jar 包,用 spark-submit 提交:
spark-submit \ --master local[*] \ --name VideoRecommendTask \ /path/to/recommend_task.py如果是 Spark 集群模式,可以指定 yarn 或 spark:// 地址,并根据实际情况配置 executor 内存:
spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 2g \ --num-executors 4 \ /path/to/recommend_task.py6.3 Django Web 服务启动
数据库配置在settings.py中,以 MySQL 为例:
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,能看到推荐系统首页,说明部署成功。
6.4 一键启动脚本
给毕设项目配一个启动脚本会显得工程化更完善。下面是一个简单的 Shell 脚本示例:
#!/bin/bash # 启动 HDFS $HADOOP_HOME/sbin/start-dfs.sh # 启动 Spark(如果使用 Standalone 模式) $SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://localhost:7077 # 启动 Django cd /path/to/project source venv/bin/activate nohup python manage.py runserver 0.0.0.0:8000 > django.log 2>&1 & echo "System started. Visit http://127.0.0.1:8000"7. 功能测试与效果验证
7.1 用户模块测试
测试目的:验证用户注册、登录、个人信息修改功能是否正常。
操作步骤:
- 打开首页,点击注册。
- 输入用户名、密码、邮箱,提交。
- 使用注册账号登录。
- 进入个人中心,修改头像或昵称。
预期结果:
- 注册成功后跳转到登录页。
- 登录成功后首页显示当前用户名。
- 个人资料修改后能正常保存。
常见失败原因:
- MySQL 数据库没有创建,Django 连接失败。
- 密码加密方式配置错误,登录校验不通过。
7.2 视频管理测试
测试目的:验证视频上传、分类展示、视频详情页功能。
操作步骤:
- 管理员后台添加视频信息,包括标题、封面、视频地址、分类、标签。
- 前端首页按分类展示视频列表。
- 点击视频进入详情页,能正常播放并展示点赞数、收藏数。
预期结果:
- 后台新增视频后,前端列表能同步展示。
- 视频详情页能正确显示视频信息和行为按钮。
7.3 推荐结果测试
测试目的:验证 Spark 离线任务生成的推荐列表是否能正常展示。
操作步骤:
- 用两个以上测试账号,在系统中分别观看、点赞不同类型的视频。
- 执行 Spark 离线推荐任务,生成推荐结果写入 MySQL。
- 分别使用不同账号登录,查看首页推荐列表。
预期结果:
- 不同账号看到的首页推荐视频有明显差异,倾向于各自感兴趣的分类。
- 推荐列表不包含当前用户已经看过的视频。
- 新注册用户能看到热门视频或最新视频,而不是空白。
判断标准:推荐结果合理、有差异、无重复推荐。
7.4 后台管理测试
测试目的:验证管理员后台的用户管理、视频审核、数据统计功能。
操作步骤:
- 使用超级管理员账号登录 Django admin 后台。
- 查看用户列表,测试禁用/启用用户。
- 查看视频列表,测试下线违规视频。
- 查看行为统计图表。
预期结果:
- 后台操作能同步影响前端展示。
- 被禁用的用户无法登录。
- 被下线的视频不再出现在推荐列表。
8. 接口 API 与批量任务设计
8.1 Django 接口示例
推荐系统的 Web 端主要接口包括:
| 接口路径 | 方法 | 功能 |
|---|---|---|
| /api/register | POST | 用户注册 |
| /api/login | POST | 用户登录 |
| /api/video/list | GET | 视频列表 |
| /api/video/detail | GET | 视频详情 |
| /api/behavior/report | POST | 行为上报 |
| /api/recommend/list | GET | 个性化推荐列表 |
推荐列表接口示例:
from django.http import JsonResponse from .models import RecommendResult, Video def recommend_list(request): user_id = request.GET.get("user_id") if not user_id: return JsonResponse({"code": 1, "message": "缺少用户ID"}) recommend_videos = ( RecommendResult.objects .filter(user_id=user_id) .order_by("-score")[:20] ) video_list = [] for item in recommend_videos: video = Video.objects.get(id=item.video_id) video_list.append({ "video_id": video.id, "title": video.title, "cover_url": video.cover_url, "video_url": video.video_url, "category": video.category, "score": round(item.score, 4) }) return JsonResponse({"code": 0, "data": video_list})Python 请求接口的测试代码:
import requests url = "http://127.0.0.1:8000/api/recommend/list" params = {"user_id": 1} response = requests.get(url, params=params, timeout=10) print(response.status_code) print(response.json())8.2 Spark 批量推荐任务
Spark 离线任务是批量生产推荐结果的核心。任务设计一般包括四步:
- 读取用户行为日志。
- 数据清洗,过滤无效记录。
- 计算用户-物品评分矩阵和相似度矩阵。
- 为每个用户生成 Top-N 推荐列表,写回 MySQL。
这里要注意:推荐任务要和 Web 服务解耦。Web 端只负责展示推荐结果,不负责计算推荐结果。这样即使推荐任务跑几个小时,也不影响网站的访问。
批量任务建议配合定时调度工具使用,Linux 下可以直接用 crontab:
# 每天凌晨 2 点执行推荐任务 0 2 * * * /usr/local/spark/bin/spark-submit --master local[*] /path/to/recommend_task.py >> /path/to/logs/spark_task.log 2>&19. 资源占用与性能观察
9.1 Spark 任务资源观察
Spark 自带 Web UI,默认端口是4040(任务运行期间可用)。启动任务后,浏览器访问http://localhost:4040,可以看到:
- Stage 数量和每个 Stage 的耗时。
- Executor 的数量、内存使用量、任务并行度。
- Shuffle 读写数据量。
- 任务失败和重试的次数。
如果发现 Shuffle 数据量过大,通常是数据倾斜或分区不合理导致的。可以在 Spark 代码中检查 join 操作前是否做了合适的预处理,或者调整分区数。
9.2 Django 服务资源观察
Django 开发服务器是单进程模型,部署到生产环境时建议使用 Gunicorn 或 uWSGI。观察资源占用可以用:
top -p $(pgrep -f "manage.py runserver")如果网站并发量上来了,Django 请求响应变慢,优先检查:
- MySQL 慢查询日志。
- 推荐结果表是否有索引。
- 页面是否发起了大量重复 Ajax 请求。
9.3 数据量增大时的性能瓶颈
当用户量和视频量增大时,有几个关键瓶颈:
| 瓶颈点 | 原因 | 优化方向 |
|---|---|---|
| 相似度计算变慢 | 用户-物品矩阵稀疏,但笛卡尔积仍然巨大 | 先用热门物品粗筛,再做精确相似度计算 |
| MySQL 查询变慢 | 行为表数据量膨胀 | 按月分表、加索引、定期归档历史数据 |
| HDFS 小文件过多 | 每次行为上报都写一个文件 | 使用 Spark Streaming 或定时合并小文件 |
| 推荐结果实时性差 | 离线任务每天跑一次 | 增加用户最近行为加权,或引入实时流计算 |
毕设阶段做到离线推荐已经合格,性能优化可以作为论文里的“进一步工作”。
10. 常见问题与排查方法
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| Hadoop 启动失败 | NameNode 未格式化或 tmp 目录损坏 | 执行jps查看进程,查看$HADOOP_HOME/logs/hadoop-*.log | 重新hdfs namenode -format,注意先备份数据 |
| Spark 任务连接不到 HDFS | HDFS 未启动或端口不正确 | 执行hdfs dfs -ls /测试连通性 | 启动start-dfs.sh,检查core-site.xml中 fs.defaultFS |
| Spark 任务执行报内存溢出 | executor 内存配置过小 | 查看 Spark UI 中 Executor 的 GC 时间和内存曲线 | 调大--executor-memory,降低任务并行度 |
| Django 页面打不开 | Django 服务未启动或端口被占用 | 执行netstat -tlnp | grep 8000 | pkill -f runserver后重新启动 |
| 推荐结果为空 | Spark 任务未执行,或推荐结果表没有数据 | 查看 Spark 任务日志,检查recommend_result表 | 确认用户有行为数据,确认行为数据格式正确 |
| 前端跨域报错 | 前端和服务端口不一致 | 查看浏览器开发者工具的 Network 面板 | 配置 Django CORS,或统一使用同源访问 |
| 中文乱码 | MySQL 字符集配置不正确 | 查看show variables like '%character%' | 建库时设置为 utf8mb4 |
11. 最佳实践与使用建议
11.1 毕设开发节奏
建议把项目分成三个阶段推进:
- 第一个阶段:数据层和算法层。先跑通 Hadoop、Spark、MySQL,把行为日志采集和最简单的热度推荐做出来,保证数据能从网页流到数据库再流到 Spark。
- 第二个阶段:推荐算法升级。在热度推荐的基础上加入协同过滤,对比两种方案的推荐效果差异。
- 第三个阶段:Web 功能完善。完善用户模块、视频模块、后台管理、前端页面展示,并配套文档报告。
11.2 推荐效果验证方法
毕设答辩时,最怕被问到“你的推荐效果怎么证明”。建议准备一组对比数据:
- 用户 A 只看了科技类视频,推荐列表是否以科技类为主。
- 用户 B 只看了美食类视频,推荐列表是否以美食类为主。
- 新注册用户是否能看到热门视频。
- 使用相同数据,热度推荐和协同过滤推荐的结果差异。
把这些对比结果截图放进论文的“实验分析”章节,比单纯说“系统运行正常”有说服力得多。
11.3 代码与文档管理
- 源码使用 Git 管理,从第一天开始就提交,不要等到最后一起提交。
- 文档报告里要包含系统架构图、数据库 ER 图、算法流程图、核心代码说明、测试用例。
- 部署环境记录成文档,方便换电脑或换服务器后快速恢复环境。
11.4 数据合规与安全提醒
推荐系统会收集用户行为数据,毕设项目虽然以学习为主,但也要注意数据合规底线:
- 不要使用真实用户的大规模隐私数据,测试数据用自己造的模拟数据。
- 论文和演示中涉及用户信息时,使用脱敏后的匿名数据。
- 项目展示时不要包含真实的账号密码、数据库口令。
- 如果是后续商用,必须获得用户授权并符合相关法规要求。
12. 总结与下一步
这个项目最值得尝试的点,是它能让你在一个 Python 项目里同时接触 Hadoop、Spark、Django 三个技术栈。推荐算法本身并不难,难的是把“日志采集 → HDFS 存储 → Spark 计算 → MySQL 结果 → 前端展示”整条链路跑通。一旦跑通,你对离线推荐系统的理解就会比只看理论扎实很多。
建议拿到项目后,最先验证三件事:第一,Hadoop 和 Spark 能不能正常启动;第二,Django 能否连上 MySQL 并完成注册登录;第三,用两个测试账号制造不同的行为数据,跑一次 Spark 推荐任务,看推荐列表是否有差异。这三个点通了,项目的主体流程就没有大问题。
最容易踩的坑也在环境层面:Hadoop 版本和 Spark 版本不兼容、JDK 版本不匹配、MySQL 字符集没配好导致写入中文乱码。这些坑不复杂,但排查起来很耗时间,建议每一步都写清楚记录,方便还原和梳理。后续如果想继续扩展,可以尝试接入 Redis 做实时热门榜单,或者用 Flask 写出行为上报接口,再把协同过滤升级成矩阵分解或深度学习模型,这些方向都能让项目在毕设基础上再往上走一步。