Scrapy+Spark+Kafka+Spring Boot构建电影推荐系统全解析
简介一套面向毕业设计、课程设计场景的电影推荐系统完整源码包涵盖Spark推荐算法、Spring Boot后端与微信小程序前端适合Java、大数据方向学生进行项目实战。压缩包共80个文件总大小16.15MB核心代码包括44个Java后端工程文件、20个Python爬虫与数据处理脚本、7个Scala推荐算法文件另含pom.xml、配置文件和项目说明文档各模块分层清晰可直接导入运行。内容不限于框架代码还附有《基于多模型融合策略的电影推荐系统设计与实现》PDF论文以及针对豆瓣电影数据的爬虫采集、评论解析、用户评分转换等脚本可辅助理解从数据获取、离线推荐到实时流推荐的全链路实现。目前已有167人学习或下载源码经过测试按说明配置环境即可快速启动尤其适合需要完整项目参考的毕业设计、课程设计或工程实训环节。1. 为什么电影推荐系统要选SparkSpring 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( fhttps://movie.douban.com/top250?start{page}, callbackself.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_DELAY2表示每个请求间隔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字段如下字段类型说明userIdint用户映射IDmovieIdint电影映射IDratingdouble归一化评分 1.0~5.0timestamplong评分时间戳为什么要映射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(indexindex, 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, headerTrue, inferSchemaTrue) als ALS( userColuserId, itemColmovieId, ratingColrating, rank20, maxIter10, regParam0.1, coldStartStrategydrop ) model als.fit(ratings) model.save(model/als_model)参数设置要解释一下rank20决定了用户向量和物品向量的维度。维度越大模型表达能力越强但会在小数据集上过拟合我测试过10到30之间效果差异不大。maxIter10是ALS迭代次数超过15次基本不收敛。regParam0.1是L2正则化系数防止某些热门电影因子过大。coldStartStrategydrop非常关键它让ALS在预测时遇到新用户或新物品直接丢弃结果而不是生成NaN否则后续融合评分会全部变成空值。离线推荐会周期性执行生成每个用户Top N的电影列表写入HBase或MySQL供后端读取。3.2 实时推荐基于Kafka Stream的流处理实时推荐部分用了Kafka Stream响应速度比离线快一个量级。整体流程是前端上报用户行为评分、点击到Kafka的user_actiontopickafkastream消费者消费这些事件根据电影相似度矩阵实时生成新的推荐列表再写回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写KStreamString, 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读取该电影的相似电影列表 ListInteger 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 Bootpom.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查询返回该用户个性化Top10search接口做分页搜索小程序端下拉加载用。注意这里没有把逻辑写在Controller里而是通过MovieService解耦方便后面加缓存和降级。4.2 整合ES查询使用Spring Data ElasticsearchMovieService内部用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)); SearchHitsMovie 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打包运行用Mavenmvn clean package -DskipTests java -jar target/movie-recommender-0.0.1-SNAPSHOT.jar --server.port8080这个配置里最坑的是版本匹配。Spring Boot 2.5.x对应spring-data-elasticsearch4.2.x如果换成Boot 2.7.x还沿用旧的ES配置启动时会报NoSuchBeanDefinitionException。我的做法是降到2.5.x或者用RestHighLevelClient手动配置避免自动装配的坑。Kafka配置里auto-offset-resetlatest表示只消费新消息重启时不会重复处理历史日志但要注意如果前端上报行为是在你启动之前数据就会丢。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-ytrue scrolltolowerloadMore view v-for(item, index) in movieList :keyitem.movieId image :srcitem.posterUrl modeaspectFill lazy-loadtrue/image view{{ item.title }}/view view评分: {{ item.rating }}/view /view /scroll-viewsrolltolower触底加载下一页对应JS里的page 1请求。lazy-loadtrue是图片懒加载在弱网环境下特别有用不然一屏加载20张海报会卡顿。5.3 易错点修改刚进入的加载页面很多同学遇到打开小程序白屏或一直转圈调试步骤分三步先看Network里的请求是否发出去再看后端日志有没有报错最后看ES里查询的数据是不是空。最常见的坑是后端接口写的是http://localhost:8080在开发者工具里能通真机上一律不通必须换成局域网IP或线上域名。另一个坑是index.json自定义了导航栏导致顶部内容被刘海屏遮挡需要计算状态栏高度const { statusBarHeight } uni.getSystemInfoSync() this.navBarHeight statusBarHeight 44最后如果你要给小程序加列表点击进入详情页不要用navigator标签硬编码用uni-simple-router管理页面栈配合onPullDownRefresh刷新推荐列表体验会好很多。本文还有配套的精品资源点击获取