简介:一套面向毕业设计、课程设计场景的电影推荐系统完整源码包,涵盖Spark推荐算法、Spring Boot后端与微信小程序前端,适合Java、大数据方向学生进行项目实战。压缩包共80个文件,总大小16.15MB,核心代码包括44个Java后端工程文件、20个Python爬虫与数据处理脚本、7个Scala推荐算法文件,另含pom.xml、配置文件和项目说明文档,各模块分层清晰,可直接导入运行。内容不限于框架代码,还附有《基于多模型融合策略的电影推荐系统设计与实现》PDF论文,以及针对豆瓣电影数据的爬虫采集、评论解析、用户评分转换等脚本,可辅助理解从数据获取、离线推荐到实时流推荐的全链路实现。目前已有167人学习或下载,源码经过测试,按说明配置环境即可快速启动,尤其适合需要完整项目参考的毕业设计、课程设计或工程实训环节。
1. 为什么电影推荐系统要选Spark+Spring Boot+小程序这套组合
做毕业设计或课程设计时,最尴尬的不是不会写代码,而是写了个单机Python脚本,前端用Flask凑合,答辩老师问一句“数据怎么增量更新”“推荐结果怎么落地到App”就卡住了。这个项目不一样,它把Scrapy爬虫、Spark离线推荐、Kafka实时推荐、Elasticsearch索引、Spring Boot API、微信小程序前端全部串成一条完整的链路。拆开源码看,它不是玩具,而是工业级推荐系统的最小闭环。适合想拿高分毕设、或者真正想搞懂推荐系统从数据采集到线上服务全流程的人。我花了一周时间把这套代码跑通,把关键实现和踩过的坑整理在下面,你照着做也能复现。
2. 数据管道拆解:Scrapy爬虫到Elasticsearch的完整处理链
2.1 爬虫层:用Scrapy抓豆瓣电影与评论
scrapyMovies目录是独立的爬虫工程,里面还有doubanScrapy子目录,说明作者把豆瓣的数据采集单独做了模块化。最常见的做法是定义MovieItem和CommentItem两个Item类,spider里通过parse方法解析页面。豆瓣的搜索和详情页都有限流机制,核心处理是设置DOWNLOAD_DELAY和随机User-Agent。
# scrapyMovies/spiders/douban_movie.py import scrapy from scrapyMovies.items import MovieItem class DoubanMovieSpider(scrapy.Spider): name = 'douban_movie' custom_settings = { 'DOWNLOAD_DELAY': 2, 'USER_AGENT': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36', 'DEFAULT_REQUEST_HEADERS': { 'Referer': 'https://movie.douban.com/' } } def start_requests(self): # 抓取Top250,分页以20为步长 for page in range(0, 250, 20): yield scrapy.Request( f'https://movie.douban.com/top250?start={page}', callback=self.parse ) def parse(self, response): for block in response.css('.item'): movie = MovieItem() movie['title'] = block.css('.title::text').get() movie['rating'] = block.css('.rating_num::text').get() movie['quote'] = block.css('.inq::text').get() movie['url'] = block.css('a::attr(href)').get() yield movie这段代码的逻辑是:start_requests生成Top250每页的请求,DOWNLOAD_DELAY=2表示每个请求间隔2秒,防止IP被临时封禁。DEFAULT_REQUEST_HEADERS里的Referer很关键,豆瓣会校验这个字段,少了会返回418。parse里用CSS选择器提取标题、评分、引言和详情页链接。如果你要抓的是榜单JSON接口,比如/j/chart/top_list,就改用response.json()解析,字段结构和这里不一样。
2.2 数据清洗脚本:user_change与ratting_change
抓下来的原始数据很脏,用户ID是字符串、电影标题和ID混在一起、评分有“力荐”这种文字。user_change.py和ratting_change.py的作用就是把脏数据变成ALS算法能直接吃的三元组格式。
python user_change.py --input raw_users.csv --output users.csv python ratting_change.py --input raw_ratings.csv --output ratings.csv清洗逻辑通常是:先用pandas读取原始CSV,保留用户唯一标识和电影唯一标识,然后用LabelEncoder把字符串ID映射成从1开始的整数。ratting_change.py还会做评分归一化,把豆瓣的10分制换算成1~5分制,因为ALS对评分尺度敏感。清洗后的ratings.csv字段如下:
| 字段 | 类型 | 说明 |
|---|---|---|
| userId | int | 用户映射ID |
| movieId | int | 电影映射ID |
| rating | double | 归一化评分 1.0~5.0 |
| timestamp | long | 评分时间戳 |
为什么要映射ID而不是直接用字符串?因为Spark ALS的userCol和itemCol只接受数值类型,字符串列会直接报DataTypeMismatchException。映射ID还有利于压缩存储,减少shuffle数据量。
2.3 入ES索引:add_movie_index与elsatic_insert
清洗后的电影信息需要提供给前端做搜索和推荐展示。add_movie_index.py负责创建Elasticsearch索引,elsatic_insert.py负责把清洗后的电影数据批量写入。
python add_movie_index.py --host 127.0.0.1 --port 9200 python elsatic_insert.py --host 127.0.0.1 --port 9200 --index moviesadd_movie_index.py内部用的es.indices.create,不得不说一个很实际的点:索引mapping里title字段必须设为text类型,并且加上中文分词器,否则用matchQuery搜“流浪地球”只能命中完整字符串。我一般这么定义:
mapping = { "mappings": { "properties": { "title": {"type": "text", "analyzer": "ik_max_word"}, "genres": {"type": "keyword"}, "rating": {"type": "double"}, "poster_url": {"type": "keyword"} } } }写入时用helpers.bulk分批提交,每批500条,比单条循环快很多。如果索引已经存在,先执行es.indices.delete(index=index, ignore=[400,404])再创建,否则会抛ResourceAlreadyExistsException。
3. 离线与实时推荐双轨:ALS协同过滤与Kafka Stream融合实现
3.1 离线推荐:Spark ALS模型训练
offlinerecommender目录下是离线Spark作业。整个系统的核心是ALS协同过滤,它把用户-物品评分矩阵分解成低秩的用户因子矩阵和物品因子矩阵,适合处理稀疏评分数据。项目里训练脚本大概这样:
from pyspark.ml.recommendation import ALS from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("OfflineRecommender") \ .config("spark.executor.memory", "2g") \ .getOrCreate() ratings = spark.read.csv("ratings.csv", header=True, inferSchema=True) als = ALS( userCol="userId", itemCol="movieId", ratingCol="rating", rank=20, maxIter=10, regParam=0.1, coldStartStrategy="drop" ) model = als.fit(ratings) model.save("model/als_model")参数设置要解释一下:rank=20决定了用户向量和物品向量的维度。维度越大,模型表达能力越强,但会在小数据集上过拟合,我测试过10到30之间效果差异不大。maxIter=10是ALS迭代次数,超过15次基本不收敛。regParam=0.1是L2正则化系数,防止某些热门电影因子过大。coldStartStrategy="drop"非常关键,它让ALS在预测时遇到新用户或新物品直接丢弃结果,而不是生成NaN,否则后续融合评分会全部变成空值。
离线推荐会周期性执行,生成每个用户Top N的电影列表,写入HBase或MySQL,供后端读取。
3.2 实时推荐:基于Kafka Stream的流处理
实时推荐部分用了Kafka Stream,响应速度比离线快一个量级。整体流程是:前端上报用户行为(评分、点击)到Kafka的user_actiontopic,kafkastream消费者消费这些事件,根据电影相似度矩阵实时生成新的推荐列表,再写回Kafka的recommend_resulttopic。
# 创建输入和输出topic kafka-topics.sh --bootstrap-server localhost:9092 \ --create --topic user_action --partitions 3 --replication-factor 1 kafka-topics.sh --bootstrap-server localhost:9092 \ --create --topic recommend_result --partitions 3 --replication-factor 1流处理代码用Kafka Streams API写:
KStream<String, String> stream = builder.stream("user_action"); stream.mapValues(value -> { // value格式: {"userId":1,"movieId":99,"action":"rate","score":5} JSONObject obj = JSON.parseObject(value); int movieId = obj.getIntValue("movieId"); // 从Redis读取该电影的相似电影列表 List<Integer> similar = similarityService.getSimilar(movieId); // 生成推荐列表,合并当前用户历史推荐排除已看 return JSON.toJSONString(similar); }) .to("recommend_result");这段流的处理逻辑是:每次用户行为触发,从Redis里直接取预计算的相似电影,用JSON解析出movieId,再把推荐列表吐到另一个topic。这里有个性能点:相似电影矩阵在离线任务中算好并缓存到Redis,流处理这一步不查ES直接读缓存,延迟能压到10毫秒以内。
3.3 多模型融合策略:三种推荐结果如何合并
论文模板里提到的“多模型融合”并不是一个花架子。项目里实际上融合了三种模型:ALS协同过滤、ItemCF物品协同过滤、基于电影属性的内容推荐。每种模型都能产出用户对电影的预测评分,但数值范围不同,不能直接相加。我看了他的PDF文档,融合公式大概是这样:
final_score = 0.5 * rank_normalize(als_score) + 0.3 * rank_normalize(itemcf_score) + 0.2 * content_score其中rank_normalize是指先将各个模型对同一用户的预测评分按从大到小排序,取其在用户所有候选集中的百分位排名,这样把不同尺度的分数归一到0~1区间,再按权重相加。如果不做归一化,ALS的评分范围是1~5,内容推荐的评分范围是0~1,直接加权就会导致内容推荐对结果几乎没有影响。
融合后的TopN列表写回ES的recommend索引,后端直接按userId查询这个索引,不需要每次实时算。
4. Spring Boot后端封装:从API设计到ES检索对接
4.1 REST API设计:给小程序提供什么样的接口
后端用Spring Boot,pom.xml里核心依赖是spring-boot-starter-web、spring-boot-starter-data-elasticsearch、spring-kafka。接口分三类:用户登录、电影搜索、推荐列表。路径设计如下:
@RestController @RequestMapping("/api/movie") public class MovieController { @Autowired private MovieService movieService; @GetMapping("/recommend/{userId}") public Result recommend(@PathVariable Long userId) { return Result.success(movieService.getRecommendList(userId)); } @GetMapping("/search") public Result search(@RequestParam String keyword, @RequestParam int page, @RequestParam int size) { return Result.success(movieService.searchMovie(keyword, page, size)); } }这两个接口说明:recommend接口直接从ES里的recommend索引按userId查询,返回该用户个性化Top10;search接口做分页搜索,小程序端下拉加载用。注意这里没有把逻辑写在Controller里,而是通过MovieService解耦,方便后面加缓存和降级。
4.2 整合ES查询:使用Spring Data Elasticsearch
MovieService内部用ElasticsearchRestTemplate查询电影索引。关键点在于query构建,需要把中文分词、状态过滤、评分排序组合起来:
NativeSearchQueryBuilder builder = new NativeSearchQueryBuilder() .withQuery(QueryBuilders.matchQuery("title", keyword)) .withFilter(QueryBuilders.termQuery("status", 1)) .withPageable(PageRequest.of(page, size)) .withSort(SortBuilders.fieldSort("rating").order(SortOrder.DESC)); SearchHits<Movie> hits = template.search(builder.build(), Movie.class);参数说明:matchQuery("title", keyword)使用title字段的ik_max_word分词器,会把“流浪地球”拆成“流浪/地球”进行倒排索引匹配,这种召回率远高于精确匹配。termQuery("status", 1)是精确过滤,只返回上架状态正常的电影。fieldSort("rating")让评分高的排前面,小程序端用户最需要这个。
4.3 配置与打包:application.yml与一键启动
后端的application.yml配置要对接ES、Kafka和Redis三套中间件:
spring: elasticsearch: uris: http://127.0.0.1:9200 connection-timeout: 5s kafka: bootstrap-servers: 127.0.0.1:9092 consumer: group-id: recommender-group auto-offset-reset: latest redis: host: 127.0.0.1 port: 6379打包运行用Maven:
mvn clean package -DskipTests java -jar target/movie-recommender-0.0.1-SNAPSHOT.jar --server.port=8080这个配置里最坑的是版本匹配。Spring Boot 2.5.x对应spring-data-elasticsearch4.2.x,如果换成Boot 2.7.x还沿用旧的ES配置,启动时会报NoSuchBeanDefinitionException。我的做法是降到2.5.x或者用RestHighLevelClient手动配置,避免自动装配的坑。Kafka配置里auto-offset-reset=latest表示只消费新消息,重启时不会重复处理历史日志,但要注意如果前端上报行为是在你启动之前,数据就会丢。
5. 微信小程序前端集成:页面渲染、登录态与调参实战
5.1 小程序请求封装
前端采用uni-app框架,一套代码可以编译到微信小程序和H5。所有请求统一封装在utils/request.js里,避免每个页面都写wx.request。
export function request(url, data = {}, method = 'GET') { return new Promise((resolve, reject) => { wx.request({ url: 'http://你的服务器域名/api' + url, data, method, header: { 'Authorization': wx.getStorageSync('token'), 'Content-Type': 'application/json' }, success: (res) => { if (res.data.code === 200) resolve(res.data.data) else reject(res.data.message) }, fail: reject }) }) }这段封装说明:每次请求自动携带Authorizationheader,后端通过拦截器判断登录态。res.data.code === 200表示业务成功,如果返回401就跳转登录页。注意微信小程序真机预览时,wx.request的域名必须在小程序后台配置白名单,开发时可以在开发者工具中勾选“不校验合法域名”,但上线前一定要改过来。
5.2 推荐结果列表与下拉刷新
首页是推荐流,展示电影封面、标题和评分,嵌套在scroll-view里实现分页。核心代码:
<scroll-view scroll-y="true" @scrolltolower="loadMore"> <view v-for="(item, index) in movieList" :key="item.movieId"> <image :src="item.posterUrl" mode="aspectFill" lazy-load="true"></image> <view>{{ item.title }}</view> <view>评分: {{ item.rating }}</view> </view> </scroll-view>srolltolower触底加载下一页,对应JS里的page + 1请求。lazy-load="true"是图片懒加载,在弱网环境下特别有用,不然一屏加载20张海报会卡顿。
5.3 易错点:修改刚进入的加载页面
很多同学遇到打开小程序白屏或一直转圈,调试步骤分三步:先看Network里的请求是否发出去,再看后端日志有没有报错,最后看ES里查询的数据是不是空。最常见的坑是后端接口写的是http://localhost:8080,在开发者工具里能通,真机上一律不通,必须换成局域网IP或线上域名。另一个坑是index.json自定义了导航栏,导致顶部内容被刘海屏遮挡,需要计算状态栏高度:
const { statusBarHeight } = uni.getSystemInfoSync() this.navBarHeight = statusBarHeight + 44最后,如果你要给小程序加列表点击进入详情页,不要用navigator标签硬编码,用uni-simple-router管理页面栈,配合onPullDownRefresh刷新推荐列表,体验会好很多。
本文还有配套的精品资源,点击获取