尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spark+Hadoop+Django:从零搭建个性化短视频推荐系统

Spark+Hadoop+Django:从零搭建个性化短视频推荐系统 基于 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_k10): 基于用户的协同过滤推荐 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(keylambda x: x[1], reverseTrue) 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(), keylambda x: x[1], reverseTrue) 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(blankTrue, nullTrue, verbose_name头像) created_at models.DateTimeField(auto_now_addTrue, verbose_name注册时间) class Meta: db_table user verbose_name 用户 class Video(models.Model): title models.CharField(max_length200, verbose_name标题) cover_url models.URLField(verbose_name封面地址) video_url models.URLField(verbose_name视频地址) category models.CharField(max_length50, verbose_name分类) tags models.CharField(max_length200, verbose_name标签) uploader models.ForeignKey(User, on_deletemodels.CASCADE, verbose_name上传者) created_at models.DateTimeField(auto_now_addTrue, 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_deletemodels.CASCADE, verbose_name用户) video models.ForeignKey(Video, on_deletemodels.CASCADE, verbose_name视频) behavior_type models.CharField(max_length20, choicesBEHAVIOR_CHOICES, verbose_name行为类型) created_at models.DateTimeField(auto_now_addTrue, 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_iduser_id, video_idvideo_id, behavior_typebehavior_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.8Hadoop 和 Spark 运行依赖必须安装并配置 JAVA_HOMEHadoop分布式存储HDFS开发环境可单机伪分布式部署Spark离线计算引擎依赖 Hadoop负责跑推荐任务MySQL业务数据库存用户、视频、行为、推荐结果Python 3Django Web 开发建议用虚拟环境管理依赖DjangoWeb 框架安装 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 创建独立环境避免污染系统 Pythonpython3 -m venv venv source venv/bin/activate pip install django pymysql requests6. 安装部署与启动方式6.1 Hadoop 伪分布式部署开发学习阶段不需要搭建真实集群单机伪分布式模式足够跑通流程。先修改$HADOOP_HOME/etc/hadoop/core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/usr/local/hadoop/tmp/value /property /configuration再修改hdfs-site.xml设置副本数为 1configuration property namedfs.replication/name value1/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 21 echo System started. Visit http://127.0.0.1:80007. 功能测试与效果验证7.1 用户模块测试测试目的验证用户注册、登录、个人信息修改功能是否正常。操作步骤打开首页点击注册。输入用户名、密码、邮箱提交。使用注册账号登录。进入个人中心修改头像或昵称。预期结果注册成功后跳转到登录页。登录成功后首页显示当前用户名。个人资料修改后能正常保存。常见失败原因MySQL 数据库没有创建Django 连接失败。密码加密方式配置错误登录校验不通过。7.2 视频管理测试测试目的验证视频上传、分类展示、视频详情页功能。操作步骤管理员后台添加视频信息包括标题、封面、视频地址、分类、标签。前端首页按分类展示视频列表。点击视频进入详情页能正常播放并展示点赞数、收藏数。预期结果后台新增视频后前端列表能同步展示。视频详情页能正确显示视频信息和行为按钮。7.3 推荐结果测试测试目的验证 Spark 离线任务生成的推荐列表是否能正常展示。操作步骤用两个以上测试账号在系统中分别观看、点赞不同类型的视频。执行 Spark 离线推荐任务生成推荐结果写入 MySQL。分别使用不同账号登录查看首页推荐列表。预期结果不同账号看到的首页推荐视频有明显差异倾向于各自感兴趣的分类。推荐列表不包含当前用户已经看过的视频。新注册用户能看到热门视频或最新视频而不是空白。判断标准推荐结果合理、有差异、无重复推荐。7.4 后台管理测试测试目的验证管理员后台的用户管理、视频审核、数据统计功能。操作步骤使用超级管理员账号登录 Django admin 后台。查看用户列表测试禁用/启用用户。查看视频列表测试下线违规视频。查看行为统计图表。预期结果后台操作能同步影响前端展示。被禁用的用户无法登录。被下线的视频不再出现在推荐列表。8. 接口 API 与批量任务设计8.1 Django 接口示例推荐系统的 Web 端主要接口包括接口路径方法功能/api/registerPOST用户注册/api/loginPOST用户登录/api/video/listGET视频列表/api/video/detailGET视频详情/api/behavior/reportPOST行为上报/api/recommend/listGET个性化推荐列表推荐列表接口示例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_iduser_id) .order_by(-score)[:20] ) video_list [] for item in recommend_videos: video Video.objects.get(iditem.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, paramsparams, timeout10) 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 219. 资源占用与性能观察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 任务连接不到 HDFSHDFS 未启动或端口不正确执行hdfs dfs -ls /测试连通性启动start-dfs.sh检查core-site.xml中 fs.defaultFSSpark 任务执行报内存溢出executor 内存配置过小查看 Spark UI 中 Executor 的 GC 时间和内存曲线调大--executor-memory降低任务并行度Django 页面打不开Django 服务未启动或端口被占用执行netstat -tlnp | grep 8000pkill -f runserver后重新启动推荐结果为空Spark 任务未执行或推荐结果表没有数据查看 Spark 任务日志检查recommend_result表确认用户有行为数据确认行为数据格式正确前端跨域报错前端和服务端口不一致查看浏览器开发者工具的 Network 面板配置 Django CORS或统一使用同源访问中文乱码MySQL 字符集配置不正确查看show variables like %character%建库时设置为 utf8mb411. 最佳实践与使用建议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 写出行为上报接口再把协同过滤升级成矩阵分解或深度学习模型这些方向都能让项目在毕设基础上再往上走一步。
返回列表